DataHub Dagster 集成指南:通过 Sensor 自动捕获管道元数据与表血缘
DataHub Dagster 集成指南通过 Sensor 自动捕获管道元数据与表血缘【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本指南完整讲解 DataHub 官方提供的 Dagster 集成方案如何通过acryl-datahub-dagster-plugin插件在 Dagster 中注册一个 DataHub Sensor在每次 Pipeline 运行后自动向 DataHub 上报 Pipeline/Task 元数据、运行状态以及表级血缘Table Lineage。读完本文你将掌握插件的安装与传感器注册方式、全部配置项的含义与默认值以及四种捕获表血缘的实战方案并能利用仓库中的完整示例与源码定位排查问题。集成概述能从 Dagster 提取什么根据 docs/lineage/dagster.md该集成连接器支持从 Dagster 提取三类元数据Dagster Pipeline 与 Task 元数据对应 DataHub 中的 DataFlow 与 DataJob 实体Pipeline 运行状态成功 / 失败 / 取消以及各步骤的运行统计表血缘Table Lineage通过 SQL 解析、IO Manager 与显式声明等方式提取。版本兼容性官方说明指出该集成已在Dagster 1.7.0上验证通过但这并不意味着它不兼容更早的版本。从源码看插件确实在刻意兼容旧版本例如在 datahub_sensors.py 中SensorReturnTypesUnion的导入被放在try/except中优先尝试 Dagster 1.9.1 的新路径失败则回退到旧的RawSensorEvaluationFunctionReturn在 dagster_generator.py 中JobSnapshot的导入也针对 Dagster 1.8.12 之后的变更做了兼容处理。使用 DataHub 的 Dagster SensorDagster 的Sensor传感器机制允许你在 Dagster 中发生重要事件时执行动作。DataHub 的 Dagster Sensor 会在每次 Dagster Pipeline 运行结束后自动向 DataHub 上报元数据可产出 Pipelines、Tasks 与运行结果等实体。关于 Dagster Sensor 的通用概念可参考官方 Dagster 文档。前置条件创建一个 Dagster 项目创建Definitions类或Repositories通过 UI 创建新项目时默认使用Definitions类来定义 Pipeline确保 DataHub GMS 服务已启动并可通过网络访问。安装插件在 Dagster 环境中安装 DataHub Dagster 插件包pip install acryl_datahub_dagster_plugin注意docs/lineage/dagster.md中使用的包名为acryl_datahub_dagster_plugin下划线而插件模块自带的 README.md 中使用的是 PyPI 规范写法acryl-datahub-dagster-plugin连字符二者指向同一发行包。注册 DataHub Sensor在启动 Dagster UI 之前将插件提供的 DataHub Sensor 导入到你的Definitions或Repository中。仓库中提供了完整的可运行示例 basic_setup.pyfrom dagster import Definitions from datahub.ingestion.graph.client import DatahubClientConfig from datahub_dagster_plugin.sensors.datahub_sensors import ( DatahubDagsterSourceConfig, make_datahub_sensor, ) config DatahubDagsterSourceConfig( datahub_client_configDatahubClientConfig( serverhttps://your_datahub_url/gms, tokenyour_datahub_token ), dagster_urlhttps://my-dagster-cloud.dagster.cloud, ) datahub_sensor make_datahub_sensor(configconfig) defs Definitions( sensors[datahub_sensor], )make_datahub_sensor的完整签名定义在 datahub_sensors.py它还支持name传感器名称默认为datahub_sensor、minimum_interval_seconds两次传感器评估之间的最小间隔秒数、default_status默认启停状态默认STOPPED可在 Dagit 或通过 GraphQL API 覆盖、monitored_jobs/job_selection/monitor_all_code_locations用于限定监控哪些 Job 或代码位置等可选参数。Sensor 内部结构从源码可以看到make_datahub_sensor实际会创建一组内部传感器见 DatahubSensors.initdatahub_success_sensor监听DagsterRunStatus.SUCCESSdatahub_failure_sensor监听DagsterRunStatus.FAILUREdatahub_canceled_sensor监听DagsterRunStatus.CANCELEDdatahub_multi_asset_sensor监控全部资产AssetSelection.all()用于更新资产组名称缓存。传感器在实例化时还会通过self.graph.test_connection()datahub_sensors.py测试与 DataHub 的连接若 GMS 不可达会在此处报错。配置选项DataHub Sensor 内部使用DatahubDagsterSourceConfig作为配置载体该数据类定义于 dagster_generator.py。各配置项及默认值如下表与官方文档一致配置项默认值说明datahub_client_config必填DataHub 客户端配置通过DatahubClientConfig指定server、token等dagster_url无Dagster Webserver 的 URL用于在 DataHub 中生成指向 Dagster Job / Run 的链接如https://myDagsterCloudEnvironment.dagster.cloud/prodcapture_asset_materializationTrue是否在 AssetMaterialization 事件发生时将资产键Asset Key捕获为 DataHub Datasetcapture_input_outputFalse是否从HANDLED_OUTPUT、LOADED_INPUT事件中捕获并解析输入/输出当前仅支持PathMetadataValue元数据标记为实验特性platform_instance无可选。该 recipe 产生的所有资产所属的平台实例asset_lineage_extractor无自定义资产血缘提取函数用于实现自己的血缘捕获逻辑enable_asset_query_metadata_parsingTrue是否启用从资产元数据中解析查询SQL以提取血缘connect_ops_to_opsFalse是否根据执行顺序将 Op 与 Op 相连在 DataJob 之间建立上游关系capture_dataset_from_asset_keyTrue是否从资产键捕获 Datasetemit_modeASYNC写往 DataHub 的发射模式。ASYNC默认避免每次写入都同步提交在高吞吐时降低 GMS 负载需要读后写一致或失败即报错时用SYNC_WAIT/SYNC_PRIMARYasset_keys_to_dataset_urn_converter无自定义资产键到 Dataset URN 的转换函数materialize_dependenciesFalse是否在 DataHub 中物化资产依赖为每个依赖发出datasetKeyemit_queriesFalse是否发出 Query 相关 aspectemit_assetsTrue是否发出资产相关 aspectdebug_modeFalse是否开启调试模式会输出更多日志例如打印每条已发射的 MCP注意datahub_client_config在数据类中是必填字段。不过源码中保留了向后兼容分支若未传入 config 对象会使用默认的DatahubClientConfig(serverhttp://localhost:8080, client_modeClientMode.INGESTION, datahub_componentdagster-plugin)并发出弃用警告见 datahub_sensors.py。启动并开启 Sensor启动 Dagster UI点击Overview标签页再进入Sensors标签页找到名为datahub_sensor的传感器可在make_datahub_sensor(name...)中自定义将其开关拨到开启状态。开启后DataHub Sensor 就会在每次 Pipeline 运行结束后自动向 DataHub 发射元数据。如何验证安装进入 Dagster UI前往OverviewSensors确认能看到datahub_sensor启动一个 Dagster Job然后在 daemon 日志中查找 DataHub 相关日志datahub_sensor - Emitting metadata...看到这条日志即表示传感器已正确配置正在向 DataHub 发射元数据。从源码看发射流程_emit_metadatadatahub_sensors.py是核心入口它会依次加载资产定义与 Job 快照 → 解析 Dagster 环境云端 / 分支部署 / 模块 / 仓库 / 代码位置见get_dagster_environment→ 收集运行日志并提取输入/输出 Dataset → 生成并发射 DataFlow对应 Dagster Joborchestrator固定为dagster→ 发射 Job Run对应 DataHub DataProcessInstance携带steps_succeeded、steps_failed、materializations、start_time、end_time等运行统计属性→ 为每个 Op 生成 DataJob 并发射 Op Run。其中失败 / 取消且未收集到任何血缘的运行会以generate_lineageFalse的方式发射 DataJob 实体避免用空血缘覆盖上一次成功运行留下的血缘关系见 datahub_sensors.py。捕获表血缘Dagster 集成提供了多种提取表血缘表与表之间的上下游关系的方式官方推荐根据场景选用以下一种或多种方案组合使用。方案一从 SQL 查询解析血缘第一步提取资产标识Asset Identifier在命名 Dagster Asset 时官方推荐使用如下结构key_prefix[env, platform, db_name, schema_name]这一命名约定保证插件能把 Asset 名称正确解析为 DataHub 的 Dataset URN。例如asset( key_prefix[prod, snowflake, db_name, schema_name], # fqdn 资产名用于识别平台并保证资产唯一 deps[iris_dataset], )只要按此规则命名 Asset插件就能在 Asset 与它指向的很可能已存在于 DataHub 中的数据集之间建立关联为后续血缘追踪打下基础。完整的同名示例见 iris.py。如果你采用了其他命名约定可以自定义asset_keys_to_dataset_urn_converter回调函数用任意方式基于元数据或其他生成 DataHub Dataset URN。以下是文档给出的、与上述命名约定配套的默认转换逻辑def asset_keys_to_dataset_urn_converter( self, asset_key: Sequence[str] ) - Optional[DatasetUrn]: Convert asset key to dataset urn By default, we assume the following asset key structure: key_prefix[prod, snowflake, db_name, schema_name] if len(asset_key) 3: return DatasetUrn( platformasset_key[1], envasset_key[0], name..join(asset_key[2:]), ) else: return None源码中的默认实现与文档完全一致dagster_generator.pyplatform取asset_key[1]、env取asset_key[0]、数据集名取asset_key[2:]以.拼接且要求资产键长度至少为 3否则返回None无法解析时该资产不会生成下游 URN。第二步为 Asset 附加 Query 元数据DataHub 的 Dagster 集成能够通过分析 Software Defined Asset 实际执行的 SQL 查询自动检测其数据集输入与输出。启用方式很简单把执行过的查询以Query标签写入 Asset 元数据即可。示例asset(key_prefix[prod, snowflake, db_name, schema_name]) def my_asset_table_a(snowflake: SnowflakeResource) - MaterializeResult: query create or replace table db_name.schema_name.my_asset_table_a as ( SELECT * FROM db_name.schema_name.my_asset_table_b ); with snowflake.get_connection() as connection: with connection.cursor() as cursor: cursor.execute(query) return MaterializeResult( metadata{ Query: MetadataValue.text(query), } )在这个例子中插件会自动识别出上游血缘为db_name.schema_name.my_asset_table_b。注意事项正确的资产命名至关重要因为查询解析器会根据生成的 URN 判断查询语言。在上例中URN 中的平台为snowflake因此按 Snowflake 方言解析。该特性受配置项enable_asset_query_metadata_parsing默认True控制。底层实现parse_sql方法datahub_sensors.py调用datahub.sql_parsing.sqlglot_lineage.create_lineage_sql_parsed_result传入查询文本、平台、平台实例、环境与默认数据库返回的SqlParsingResult中的in_tables与out_tables分别构成上游与下游。此外_process_lineagedatahub_sensors.py会从资产物化日志的Query元数据中取出查询文本若同时开启了emit_queries默认False还会通过gen_query_aspect将查询作为 DataHub 的 Query 实体发出dagster_generator.py该查询实体以get_query_fingerprint(query, platform)生成的指纹作为 ID。仓库中的完整示例 iris.py 同时演示了 SnowflakeTEST_DB.public.iris_setosa依赖TEST_DB.public.iris_cleaned与 Redshiftpublic.blood_storage两种平台的 SQL 血缘捕获。方案二使用 DataHubSnowflakePandasIOManager插件提供了 Dagster 原生SnowflakePandasIOManager的增强版本DataHubSnowflakePandasIOManager。该版本会自动捕获由 IO Manager 创建的 Snowflake 资产并为 Dagster 中的资产附加 DataHub URN 和跳转链接。使用方法将SnowflakePandasIOManager直接替换为DataHubSnowflakePandasIOManager即可。增强版额外接受两个参数datahub_base_urlDataHub UI 的基础 URL用于生成指向 DataHub 中 Snowflake Dataset 的直接链接未设置则不生成链接。datahub_env生成 URN 时使用的 DataHub 环境默认PROD。示例from datahub_dagster_plugin.modules.snowflake_pandas.datahub_snowflake_pandas_io_manager import ( DataHubSnowflakePandasIOManager, ) # ... resources{ snowflake_io_manager: DataHubSnowflakePandasIOManager( databaseMY_DB, accountmy_snowflake_account, warehouseMY_WAREHOUSE, usermy_user, passwordmy_password, rolemy_role, datahub_base_urlhttp://localhost:9002, ), }底层原理DataHubSnowflakePandasIOManager继承自SnowflakePandasIOManager其create_io_manager方法构造了一个DataHubDbIoManager见 datahub_snowflake_pandas_io_manager.py。核心逻辑在 datahub_db_io_manager.py 的handle_output中在 IO Manager 正常写完数据后根据table_slice的 database / schema / table 拼装出形如urn:li:dataset:(urn:li:dataPlatform:snowflake,db.schema.table,PROD)的 URN并通过context.add_output_metadata写入datahub_urn以及可选的datahub_url元数据。这些元数据会在 AssetMaterialization 事件中被传感器读取见_get_asset_downstream_urn优先使用datahub_urn元数据其次才回退到资产键转换从而实现“IO Manager 写入哪张表血缘下游就指向哪张表”。方案三使用 Dagster Ins 和 Out 显式声明对于 Assets 和 Ops都可以通过一个与装饰函数参数对应的Ins/Out字典显式提供输入和输出并在提供输入输出时附带额外元数据。利用带元数据的 ins / out 字典即可为 Assets 和 Ops 创建数据集上游与下游依赖。Assets 示例见 assets_job.pymulti_asset( outs{ extract: AssetOut( metadata{datahub.outputs: [DatasetUrn(snowflake, tableD).urn()]} ), } ) def extract(): ... asset( ins{ extract: AssetIn( extract, metadata{datahub.inputs: [DatasetUrn(snowflake, tableC).urn()]}, ) } ) def transform(extract): ...Ops 示例见 ops_job.pyop( ins{ data: In( dagster_typePythonObjectDagsterType(list), metadata{datahub.inputs: [DatasetUrn(snowflake, tableA).urn()]}, ) }, out{ result: Out( metadata{datahub.outputs: [DatasetUrn(snowflake, tableB).urn()]} ) }, ) def transform(data): ...底层原理传感器在generate_datajob中遍历每个 Op 快照的input_def_snaps与output_def_snaps当检测到datahub.inputs/datahub.outputs元数据键时会将其中的 URN 字符串解析为DatasetUrn并加入 DataJob 的inlets/outlets见 dagster_generator.py。需要说明输出端DATAHUB_OUTPUTS的解析受connect_ops_to_ops配置控制。此外Op 之间的执行依赖会通过step_deps传递上游 Op 的output_datasets会自动成为下游 Op 的inlets这是默认行为。方案四自定义资产血缘提取逻辑asset_lineage_extractor你可以完全自定义资产血缘的捕获逻辑。asset_lineage_extractor回调接收三个参数返回一个Dict[str, DatasetLineage]返回字典的key 是 op key返回字典的value 是该 op 的上游 / 下游资产 URN 集合DatasetLineage包含inputs与outputs两个Set[DatasetUrn]。from datahub_dagster_plugin.client.dagster_generator import DagsterGenerator, DatasetLineage def asset_lineage_extractor( context: RunStatusSensorContext, dagster_generator: DagsterGenerator, graph: DataHubGraph, ) - Dict[str, DatasetLineage]: dataset_lineage: Dict[str, DatasetLineage] {} # Extracting input and output assets from the context return dataset_lineage仓库中有一个完整可运行的实现 advanced_ops_jobs.py它遍历运行的日志过滤ASSET_MATERIALIZATION、ASSET_OBSERVATION、HANDLED_OUTPUT、LOADED_INPUT四类事件对每条物化日志调用dagster_generator.emit_asset(...)将资产发射为 Dataset并把返回的 URN 加入该 step 的输出集合。最后通过DatahubDagsterSourceConfig(..., asset_lineage_extractorasset_lineage_extractor)传入配置。从源码看自定义提取器与默认日志解析是叠加关系_emit_metadata会先调用asset_lineage_extractor得到血缘再调用process_dagster_logs从日志中解析血缘最后用merge_dicts合并二者见 datahub_sensors.py。实体映射与运行原理为了更好地理解集成行为下表总结了插件在 DataHub 侧产生的实体映射依据 dagster_generator.py 的generate_dataflow/generate_datajob/emit_job_run/emit_op_run实现Dagster 概念DataHub 实体说明JobDataFloworchestratordagsterID 形如branch/module/job_name云端或module/job_name本地Op / Asset 步骤DataJob属于对应 DataFlowID 形如branch/module/op_nameJob RunDataProcessInstance携带运行统计属性成功/失败步骤数、物化次数、起止时间等与结果状态SUCCESS / FAILURE / SKIPPEDOp RunDataProcessInstance携带 step_key、attempts、起止时间等且会克隆 DataJob 的 inlets / outlets 以保留血缘Asset物化事件Dataset平台固定为dagster子类型Asset名称取资产键最后一段资产组标签 / 浏览路径asset_group:group_name标签与 BrowsePathsV2 路径配置了dagster_url时DataFlow、DataJob 与运行实例都会生成指向 Dagster UI 的回链云端环境会额外拼接branch与locations/location_name路径见job_url_generator。同时Op 实体会附带SubTypesClass(typeNames[Op])资产实体会附带SubTypesClass(typeNames[Asset])便于在 DataHub UI 中区分实体类型。调试与常见问题ConnectionError for DataHub Rest URL如果在日志中看到ConnectionError: HTTPConnectionPool(hostlocalhost, port8080)说明你的DataHub GMS 服务没有启动。请先启动 GMS本地开发通常监听localhost:8080与插件默认的DEFAULT_DATAHUB_REST_URL一致或者在DatahubClientConfig.server中显式指定正确的 GMS 地址后重启 Dagster。需要注意的是DatahubSensors.__init__中的test_connection()会在传感器构造阶段就校验连通性因此连接问题会尽早暴露。其他排查建议开启调试模式设置debug_modeTrue传感器会在日志中打印每条发射的 MCPEmitted MCP: ...便于核对发射内容检查查询解析失败若enable_asset_query_metadata_parsingTrue但未产生血缘可在日志中查找Error in parsing sql: ...来自parse_sql或Lineage not found for ...通常是资产未按key_prefix[env, platform, db_name, schema_name]命名或Query元数据未被写入MaterializeResult.metadata验证资产命名确认资产键长度 ≥ 3且asset_key[1]为 DataHub 中已注册的数据平台如snowflake、redshift。仓库导航关联文档docs/lineage/dagster.md插件包PyPI 发行版metadata-ingestion-modules/dagster-plugin/插件 READMEmetadata-ingestion-modules/dagster-plugin/README.md传感器核心实现datahub_sensors.py实体生成与配置定义dagster_generator.pySnowflake IO Manager 增强datahub_snowflake_pandas_io_manager.py 与 datahub_db_io_manager.py可运行示例basic_setup.py、iris.py、assets_job.py、ops_job.py、advanced_ops_jobs.py单元测试test_dagster.py含黄金文件 golden_test_emit_metadata_mcps.json 可对照发射的 MCP 内容【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考