OpenMetadata 采集框架解析:BaseWorkflow 的步骤抽象、状态汇总与"异常即数据"设计
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
本篇基于 OpenMetadata 采集框架(ingestion)的官方设计文档 workflow/README.md 展开,系统讲解BaseWorkflow如何用四类可组合的Step组织每一次采集执行、如何在框架层面统一接管异常与Status汇报,并结合仓库源码(base.py、step.py 等)印证其真实实现。读完后,你将能够理解 OpenMetadata 中任意一条 Ingestion Pipeline 从 YAML 到执行的生命周期,并具备自行阅读、扩展采集步骤代码的能力。
1. 为什么需要 BaseWorkflow
OpenMetadata 的采集体系(Metadata Ingestion、Profiler、Usage、Classification、Test Suite、Data Insights 等)数量众多,如果每个连接器各自管理执行流程,异常与状态处理很快就会失控。BaseWorkflow的设计目标(出自 workflow/README.md)可以概括为三点:
- 执行组织达成共识:所有执行都以"步骤(Steps)"为单位组织,形成统一的约定;
- 异常管理集中化:在单一位置集中处理所有异常管理,避免某个步骤中"未被正确捕获"的异常炸掉整次执行;
- Status 处理集中化:统一回答"处理了哪些资产、哪些失败了",并将结果回传给
IngestionPipeline实体。
在源码中,这一角色由 base.py 中的BaseWorkflow抽象基类承担(class BaseWorkflow(ABC, WorkflowStatusMixin)),它定义了所有工作流共用的契约:
- 类方法
create(config_dict):创建工作流实例的唯一入口; - 抽象方法
post_init():内部组件初始化完成后执行; - 抽象方法
execute_internal():各工作流自己的安全执行逻辑; - 抽象方法
get_failures()/workflow_steps():上报失败明细与步骤状态。
2. Steps:四类步骤及其业务命名
每个Workflow都可以用Step作为"乐高积木"来搭建。每个步骤是对"其中预期会发生哪类操作"的一个通用抽象。BaseWorkflow接受任意数量的顺序Step,每个步骤负责业务逻辑的一部分。
框架层面主要有四种步骤(定义见 step.py):
| 抽象步骤 | 职责 | 具体业务类 |
|---|---|---|
IterStep | 负责启动工作流:从外部世界读取数据,并yield需要在管道中继续处理的元素 | Source |
ReturnStep | 接受一个输入,进一步处理后返回一个输出 | Processor与Sink |
StageStep | 接受一个输入,将其"暂存(stage)"到某处(如文件);预期与BulkStep配合使用 | Step(Stage) |
BulkStep | 遍历由StageStep产出的输入,不返回任何东西 | BulkSink |
这些抽象名字偏学术化,不易想象。因此仓库中定义了基于它们的具体类,便于讨论工作流结构:
IterStep->SourceReturnStep->Processor与SinkStageStep->StepBulkStep->BulkSink
开发任一具体步骤时,只需实现其执行方法(IterStep是_iter,其余是_run):IterStep中的方法预期以yield产出结果,其余步骤则以return返回结果。
2.1 步骤的执行契约:Either 模型
四类步骤共享同一个"把异常当数据"的契约:每个步骤yield或return一个Either对象,表示处理单个元素的结果要么为right——包含预期的实体结果,要么为left——包含被抛出的异常。Either的真实定义在 models.py:
class Either(BaseModel, Generic[T]): """Any execution should return us Either an Entity of an error for us to handle""" left: Annotated[ StackTraceError | None, Field(description="Error encountered during execution", default=None), ] right: Annotated[T | None, Field(description="Correct instance of an Entity", default=None)]其中Entity即任意 PydanticBaseModel。从 step.py 的源码看,IterStep与ReturnStep的基类run方法会统一检查Either:若left非空则调用self.status.failed(...)记录失败;若right非空则调用self.status.scanned(...)记录成功。此外基类还会捕获一种特殊情形——_run/_iter返回的对象根本不是Either(会报Not an Either错误),确保契约不被静默破坏。
3. Workflows:从 set_steps 到 execute_internal
有了这些积木之后,就可以定义Workflow结构。虽然步骤在理论上可以比较自由地拼接,但 OpenMetadata 遵循几套固定的"配方"。每个Workflow通过定义自己的步骤(从Source开始,添加Processor等)并在set_steps方法中注册步骤来构建。BaseWorkflow负责处理公共逻辑:初始化metadata客户端对象、timer状态日志,以及按需把状态发送到IngestionPipeline。
3.1 示例一:Metadata Ingestion
该工作流只有两个步骤:
Source:列举来源端(Dashboards、Tables、Pipelines 等)的元数据,并翻译成 OpenMetadata 标准模型;REST Sink:接收上述实体的 Create Request,发送到 OpenMetadata 服务端。
工作流在这里的作用是把步骤组合在一起并让执行流水线化——由工作流本身决定如何把Source产出的每个元素传给Sink。对应源码是 metadata.py 中的MetadataWorkflow:
def set_steps(self): # We keep the source registered in the workflow self.source = self._get_source() sink = self._get_sink() self.steps = (sink,)其中_get_source()根据source.type动态导入连接器 Source 类并执行prepare(),_get_sink()则按sink.type加载 REST Sink 等目标端。
3.2 示例二:Profiler Ingestion
Profiler 工作流的步骤更多(文档描述为 4 步):
Source:从 OpenMetadata API 中取出需要 profiling 的表;Profiler Processor:对每张表执行指标计算并收集结果;PII Processor:拿到 profiler 结果后,使用 NLP 模型为表追加分类结果;REST Sink:把结果发送到 OpenMetadata API。
同样地,Workflow类负责把元素从Source->Profiler Processor->PII Processor->REST Sink逐站传递。
对应源码是 profiler.py 中的ProfilerWorkflow,它在set_steps中注册profiler_processor与sink两个后续步骤,并在初始化时执行连接测试(test_connection()):
def __init__(self, config): super().__init__(config) self.workflow_config.successThreshold = 80注意这里把successThreshold设为 80,即允许 Profiler 运行有不超过 20% 的失败率仍不算整体失败——这与第 4 节的状态阈值机制直接相关。从源码结构看,PII 分类在当前代码库中已实现为独立的 classification.py 工作流(ClassificationWorkflow),与 Profiler 工作流分离编排,这与文档所述"Processor 链"的设计意图一致。
3.3 通用执行流程:一个 flatMap 实现
所有 Ingestion 类工作流(metadata、lineage、usage、profiler、test suite、data insights)都继承 ingestion.py 中的IngestionWorkflow,其execute_internal展示了文档所说的"配方"落地方式(L156-L178):
def execute_internal(self): """ Pass each record from the source down the pipeline: Source -> (Processor) -> Sink or Source -> (Processor) -> Stage -> BulkSink """ for record in self.source.run(): processed_record = record for step in self.steps: # We only process the records for these Step types if processed_record is not None and isinstance(step, (Processor, Stage, Sink)): processed_record = step.run(processed_record) # Try to pick up the BulkSink and execute it, if needed bulk_sink = next((step for step in self.steps if isinstance(step, BulkSink)), None) if bulk_sink: bulk_sink.run()两条要点值得注意:
- Source 必须是迭代器:
self.source.run()返回的是生成器,记录逐个流过管道,天然支持"边读边写"的流式执行; None即中断传递:某一步把记录处理失败(返回None)后,该记录不再流向后续步骤——这正是文档所说的"每个Step控制自己的Status和异常(包裹在Either中),只把工作流下游传递真正的right结果"。
文档最后也点出了这一执行模型的本质:可以把它理解为一个flatMap实现——理论上可以继续往里拼接步骤,而不必改动框架本身。
4. Status:步骤状态如何汇总为工作流状态
Workflow掌控执行流,但最重要的部分在于状态处理与异常管理。设计约定:每个Step拥有自己的Status,记录处理了什么、失败了什么;整体Workflow的状态由各步骤状态汇总得出。
Status模型定义在 status.py,是一个 Pydantic 模型,核心字段包括:
records/record_count:扫描到的记录(只保留可打印的log_name,控制内存占用);updated_records:以PatchRequest/PatchedEntity形式更新的记录;warnings:警告列表;filtered:被过滤掉的实体及原因;failures:失败明细,类型为TruncatedStackTraceError——源码注释说明这是对StackTraceError的截断版(单字段上限 1MB),防止某些连接器产生爆量的异常负载。
Status.calculate_success()用(成功记录数 × 100)/(成功记录数 + 失败数)计算单步骤成功率;工作流层的calculate_success()(base.py)则对各步骤成功率做统计汇总,得到一个整体的成功率数值。
框架还用三处机制把状态"推出去":
- 周期汇报:
BaseWorkflow内置一个RepeatedTimer,每 30 秒(REPORTS_INTERVAL_SECONDS = 30,见 base.py)调用_report_ingestion_status(),按步骤打印 "Processed X records, updated X records, filtered X records, found X errors",并可向服务端发送实时进度更新(send_progress_update); - 结束时打印:
print_status()通过WorkflowOutputHandler输出最终各步骤的Summary(定义在 step.py 的Summary类,含 records / updated / warnings / errors / filtered 与失败明细); - 回传 Pipeline 实体:
build_ingestion_status()+set_ingestion_pipeline_status(...)把最终状态写回 OpenMetadata 服务端的IngestionPipeline实体。
5. Exceptions:try/catch 兜底 + 异常即数据
为了保证"所有异常都被捕获",框架采用双层策略。
第一层:每个 Step 的run方法都包在 try/catch 中。只有WorkflowFatalError(定义见 step.py,典型场景如 Test Connection 失败——此时继续执行毫无意义)会真正炸掉执行;其他任何异常只会被记录到Status。以下是文档给出的IterStep.run实现(与 step.py 源码一致):
def run(self) -> Iterable[Optional[Entity]]: """ Run the step and handle the status and exceptions Note that we are overwriting the default run implementation in order to create a generator with `yield`. """ try: for result in self._iter(): if result.left is not None: self.status.failed(result.left) yield None if result.right is not None: self.status.scanned(result.right) yield result.right except WorkflowFatalError as err: logger.error(f"Fatal error running step [{self}]: [{err}]") raise err except Exception as exc: error = f"Encountered exception running step [{self}]: [{exc}]" logger.warning(error) self.status.failed( StackTraceError( name="Unhandled", error=error, stack_trace=traceback.format_exc() ) )第二层:把异常当数据。各组件中"可能发生的各种异常",统一以Either.left的形式沿数据流传递(见 2.1 节)。这一约定的好处是:每个Step都保证把记录的异常写进自己的Status,因此所有错误都能在整次执行结束时被完整汇总、汇报。
代码中还有一个专门的观测点:Unhandled异常。当某段代码抛出了本应自己处理却漏掉的异常时,基类的兜底分支会以name="Unhandled"记录它。通过跟踪这些Unhandled异常,开发者可以知道哪些代码路径需要更审慎地处理未知场景。
6. 执行生命周期:execute() 与 successThreshold
理解了步骤与状态,最后看 base.py 中BaseWorkflow.execute()的完整编排,它串联了文档所述的全部机制:
- 启动计时器并上报"开始运行":
self.timer.trigger()启动周期状态汇报,并立即发送一条DISCOVERY进度更新,使运行在开始瞬间即对实时查看者可见; - 执行具体工作流:调用子类实现的
execute_internal(); - 判定部分成功:若
successThreshold <= 成功率 < 100,管道状态置为PipelineState.partialSuccess——这就是 3.2 节中 Profiler 把阈值设为 80 的落点; - 阈值内失败判定:
raise_from_status_internal()(base.py)逐个检查步骤,任何步骤存在失败且成功率低于workflowConfig.successThreshold时抛出WorkflowExecutionError,管道状态置为failed; - 兜底收尾:
finally块中依次执行close_steps()(关闭各步骤,让批量缓冲的 Sink 在close()中冲刷记录并计入状态——_steps_closed标志保证幂等)、build_ingestion_status()并回传服务端、print_status()打印摘要,最后stop()停止计时器、诊断线程与元数据客户端,确保"状态一定会被发送"这一不变量。
7. 在测试中验证这套设计
这套设计并非纸面约定,仓库中有专门的单元测试覆盖。test_base_workflow.py 使用桩步骤验证BaseWorkflow的状态与执行逻辑,例如:
class SimpleSource(WorkflowSource): """Simple Source for testing""" def _iter(self, *args, **kwargs) -> Iterable[Either]: for element in range(0, 5): yield Either(right=element)测试中刻意构造了 "Source not returning an Either" 的BrokenSource等场景,验证框架对契约违背(非Either返回值)的捕获,与第 2.1 节描述的Not an Either检查一一对应。同目录下的 test_status_mixin_progress.py、test_progress_rendering.py、test_application_workflow.py 等则分别覆盖进度上报、状态渲染与特殊工作流形态。
8. 小结与延伸阅读
BaseWorkflow用三件东西支撑起 OpenMetadata 全部采集管道:四类可组合步骤(Source/Processor/Stage/Sink/BulkSink)、以Either为载体的"异常即数据"约定、以及集中化的Status汇总与阈值判定。新增一个采集工作流时,只需继承IngestionWorkflow并实现set_steps(),公共的执行、异常与状态机制即自动生效。
关键源码索引:
- 设计文档:ingestion/src/metadata/workflow/README.md
- 工作流基类:ingestion/src/metadata/workflow/base.py
- 通用采集工作流:ingestion/src/metadata/workflow/ingestion.py
- 步骤抽象与 Either/Status:ingestion/src/metadata/ingestion/api/step.py、ingestion/src/metadata/ingestion/api/models.py、ingestion/src/metadata/ingestion/api/status.py
- 各业务工作流:metadata.py、profiler.py、classification.py、data_quality.py、usage.py、application.py
- 单元测试:ingestion/tests/unit/workflow/test_base_workflow.py
【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考