ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

OpenMetadata 采集框架解析:BaseWorkflow 的步骤抽象、状态汇总与“异常即数据“设计

OpenMetadata 采集框架解析:BaseWorkflow 的步骤抽象、状态汇总与“异常即数据“设计 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. 为什么需要 BaseWorkflowOpenMetadata 的采集体系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需要在管道中继续处理的元素SourceReturnStep接受一个输入进一步处理后返回一个输出Processor与SinkStageStep接受一个输入将其暂存stage到某处如文件预期与BulkStep配合使用StepStageBulkStep遍历由StageStep产出的输入不返回任何东西BulkSink这些抽象名字偏学术化不易想象。因此仓库中定义了基于它们的具体类便于讨论工作流结构IterStep-SourceReturnStep-Processor与SinkStageStep-StepBulkStep-BulkSink开发任一具体步骤时只需实现其执行方法IterStep是_iter其余是_runIterStep中的方法预期以yield产出结果其余步骤则以return返回结果。2.1 步骤的执行契约Either 模型四类步骤共享同一个把异常当数据的契约每个步骤yield或return一个Either对象表示处理单个元素的结果要么为right——包含预期的实体结果要么为left——包含被抛出的异常。Either的真实定义在 models.pyclass Either(BaseModel, Generic[T]): Any execution should return us Either an Entity of an error for us to handle left: Annotated[ StackTraceError | None, Field(descriptionError encountered during execution, defaultNone), ] right: Annotated[T | None, Field(descriptionCorrect instance of an Entity, defaultNone)]其中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 中的MetadataWorkflowdef 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 IngestionProfiler 工作流的步骤更多文档描述为 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-L178def 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. Exceptionstry/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(fFatal error running step [{self}]: [{err}]) raise err except Exception as exc: error fEncountered exception running step [{self}]: [{exc}] logger.warning(error) self.status.failed( StackTraceError( nameUnhandled, errorerror, stack_tracetraceback.format_exc() ) )第二层把异常当数据。各组件中可能发生的各种异常统一以Either.left的形式沿数据流传递见 2.1 节。这一约定的好处是每个Step都保证把记录的异常写进自己的Status因此所有错误都能在整次执行结束时被完整汇总、汇报。代码中还有一个专门的观测点Unhandled异常。当某段代码抛出了本应自己处理却漏掉的异常时基类的兜底分支会以nameUnhandled记录它。通过跟踪这些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(rightelement)测试中刻意构造了 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/Statusingestion/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),仅供参考
返回列表