ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

DataHub Airflow 插件实战指南:自动列级血缘、Assets 适配与自定义算子集成

DataHub Airflow 插件实战指南:自动列级血缘、Assets 适配与自定义算子集成 DataHub Airflow 插件实战指南自动列级血缘、Assets 适配与自定义算子集成【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读本文档面向使用 Apache Airflow 调度数据管道、并希望将管道与任务元数据、运行状态以及表级/列级血缘同步到 DataHub 的工程团队。文中以acryl-datahub-airflow-plugin为核心完整讲解安装配置、自动血缘提取、手动血缘标注inlets/outlets、原生 Airflow AssetsDatasets适配、自定义算子的三种接入方式、过时任务清理、GraphQL 查询以及常见问题排障并辅以仓库源码级佐证帮助读者在生产环境中快速落地并深度定制。说明本文聚焦于使用 Airflow 调度数据管道并采集血缘/元数据到 DataHub。如果你只是想用 Airflow 调度 DataHub 自身的元数据摄取任务即用 DAG 运行 ingestion recipe请参阅 metadata-ingestion/schedule_docs/airflow.md如果想在本地用 Docker 同时跑起 Airflow 与 DataHub 进行体验可参考 docker/airflow/local_airflow.md该指南已标记为未维护以当前仓库内容为准。插件能力总览DataHub Airflow 插件acryl-datahub-airflow-plugin在当前仓库中的源码位于 metadata-ingestion-modules/airflow-plugin/src/datahub_airflow_plugin其核心能力包括自动列级血缘提取覆盖 SQL 算子如MySqlOperator、PostgresOperator、SnowflakeOperator、BigQueryInsertJobOperator等以及S3FileTransformOperator等更多算子Airflow DAG 与任务的元数据同步包括属性properties、ownership属主与 tags标签任务运行信息任务成功/失败状态会被捕获展示在 DataHub 的 Runs 标签页中手动血缘标注通过在算子operator上设置inlets与outlets声明血缘原生 Airflow AssetsDatasets适配自动将 Asset/Dataset 映射为 DataHub Dataset URN。插件要求Airflow 3.0与Python 3.10。如果运行在更老的 Airflow 版本上需要固定使用旧版acryl-datahub-airflow-plugin详见文末 兼容性说明。从源码结构看插件的关键模块分工如下模块文件职责datahub_plugin.py/airflow3/datahub_listener.py插件入口与 Airflow listener监听 DAG/task 生命周期事件_config.py从airflow.cfg读取[datahub]配置并构建DatahubLineageConfigentities.py提供Dataset与Urn两类血缘实体负责生成/校验 DataHub URN_airflow_asset_adapter.py将 Airflow Asset/Dataset URI 转换为 DataHub Dataset URN含 URI scheme → 平台映射airflow3/_sql_parser_patch.py全局 patch OpenLineage 的SQLParser注入 DataHub 增强解析列级血缘_sql_parsing_common.py封装 DataHub SQL 解析器含多语句解析模式client/airflow_generator.py将 Airflow DAG/Task 转换为 DataHub DataFlow/DataJob 实体插件安装在满足 Airflow 3.0、Python 3.10 的前提下直接通过 pip 安装pip install acryl-datahub-airflow-plugin如果当前环境不满足上述版本要求请跳到 兼容性说明 选择对应版本的插件。连接配置插件通过 Airflow Connection 来定位 DataHub 服务端支持两种配置方式。方式一命令行添加 Connectionairflow connections add --conn-type datahub-rest datahub_rest_default --conn-host http://datahub-gms:8080 --conn-password optional datahub auth token--conn-type固定为datahub-rest--conn-hostDataHub GMS 的地址。自建部署通常为http://datahub-gms:8080若使用 DataHub Cloud则填写https://YOUR_PREFIX.acryl.io/gms--conn-password可选的 DataHub 认证 token开启认证时必填。方式二Airflow UI 添加 Connection在 Airflow UI 中进入Admin → Connections点击 新建连接在 Connection Type 下拉框中选择 DataHub REST Server填入对应的 Host 与 Password 即可。可选配置项airflow.cfg插件默认无需额外配置即可工作。所有可选参数都放在airflow.cfg的[datahub]段下[datahub] # Optional - additional config here. enabled True # default完整参数表如下结合 metadata-ingestion-modules/airflow-plugin/src/datahub_airflow_plugin/_config.py 中get_lineage_config()的实际读取逻辑整理名称默认值说明enabledtrue插件是否启用。conn_iddatahub_rest_defaultDataHub REST connection 的名称。源码中该字段支持逗号分隔的多个 connection ID_datahub_connection_ids会按逗号拆分多连接时自动使用DatahubCompositeHook聚合写入。clusterprodAirflow 集群名称等价于实例的envDataHub 环境。platform_instanceNone该插件产出的所有资产所属的平台实例可选。capture_ownership_infotrue是否提取 DAG 的属主信息owner 会映射为 DataHub corpuser。capture_ownership_as_groupfalse提取 DAG 属主时将 owner 视为 corpgroup用户组而非 corpuser。capture_tags_infotrue是否提取 DAG 标签为 DataHub tags。capture_executionstrue是否提取任务运行及其成功/失败状态展示在 DataHub Runs 标签页。materialize_ioletstrue是否创建或取消软删除血缘中引用的所有实体。enable_extractorstrue是否启用自动血缘提取。disable_openlineage_plugintrue是否禁用 OpenLineage 插件以避免重复处理。默认 true 时只有 DataHub listener 运行SQLParser 被 patch 为调用 DataHub 增强解析器设为 false 则两个插件并行DataHub 从 run_facets 读取增强解析结果OpenLineage 保留自己的 inputs/outputs。enable_multi_statement_sql_parsingfalse是否在单个任务内解析多条 SQL 语句解析临时表并在一次执行内合并血缘。log_levelno change[debug] 设置插件日志级别。debug_emitterfalse[debug] 为 true 时插件会记录所发出的事件log 出 emitter 内容。dag_filter_str{ allow: [.*] }以 JSON 字符串形式表示的 AllowDenyPattern用于过滤参与处理的 DAG。源码中通过AllowDenyPattern.model_validate_json()解析。enable_datajob_lineagetrue为 true 时插件会为 DataJob 发送输入/输出血缘。capture_airflow_assetstrue是否将原生 Airflow Assets/Datasets 捕获为 DataHub 血缘详见 原生 Airflow Assets/Datasets。emit_modeASYNC向 DataHub 写入时的发送模式。默认ASYNC异步不阻塞每次写入的同步提交降低高并发下的 GMS 负载如需读后写一致或失败即抛错可使用SYNC_WAIT/SYNC_PRIMARY。此外源码中还提供了以下未写入原文档参数表、但同样可在airflow.cfg中使用的细粒度开关仅在enable_extractorsTrue时生效名称默认值说明patch_sql_parsertrue是否 patch OpenLineage 的SQLParser.generate_openlineage_metadata_from_sql()以启用 DataHub 解析列级血缘。extract_athena_operatortrue是否用 DataHub 增强实现替换 Athena 算子的get_openlineage_facets_on_complete()。extract_bigquery_insert_job_operatortrue同上针对BigQueryInsertJobOperator。extract_teradata_operatortrue同上针对 Teradata 算子。render_templatestrue是否在 listener 中调用ti.render_templates()使 Jinja 模板字段的提取更准确。datajob_url_linktasksDataJob 详情链接的格式tasks/dags/{dag_id}/tasks/{task_id}或grid。注意旧版 Airflow 2 的taskinstance格式已移除若配置该值会直接抛出带迁移提示的ValueError见 _config.py。自动血缘提取插件的自动血缘能力建立在 Airflow 内置的 OpenLineage 支持 之上因此支持的算子集合是 Airflow/OpenLineage 默认支持的超集。在 Airflow 3.x 中SQL 算子是直接调用SQLParser.generate_openlineage_metadata_from_sql()而非使用 extractor。插件在导入时通过 airflow3/_sql_parser_patch.py 对该方法做全局 patch先按需调用原始解析器当disable_openlineage_pluginFalse时OpenLineage 保留自己的 inputs/outputs再调用 DataHub 增强解析器解析结果存入自定义 facetDATAHUB_SQL_PARSING_RESULT_KEY供 listener 读取。SQL 相关的 extractor 已改用 DataHub 的 SQL 血缘解析器比内置解析器更健壮并能利用 DataHub 的元数据信息生成列级血缘。受支持的算子清单SQLExecuteQueryOperator及其所有子类绝大多数 SQL 算子都继承自该类AthenaOperator与AWSAthenaOperatorBigQueryOperator与BigQueryExecuteQueryOperatorBigQueryInsertJobOperatorincubatingMySqlOperatorPostgresOperatorRedshiftSQLOperatorSnowflakeOperator与SnowflakeOperatorAsyncSqliteOperatorTeradataOperator注意Teradata 使用不含 schema 层级的两级database.table命名TrinoOperator。多语句 SQL 解析Multi-Statement SQL Parsing当一个任务执行多条 SQL 语句时例如CREATE TEMP TABLE ...; INSERT ... FROM temp_table;默认False只解析第一条语句。开启该选项后插件会将所有语句合并解析解析临时表依赖并在一次执行中合并血缘[datahub] enable_multi_statement_sql_parsing True # Default: False从 _sql_parsing_common.py 的parse_sql_with_datahub()可以看到两种模式的底层差异开启时调用create_lineage_from_sql_statements()将整个 SQL 语句列表一次性交给 DataHub 解析器解析并合并临时表关闭时调用create_lineage_sql_parsed_result()且若传入的是列表则只取第一条语句进行解析sql sql[0] if sql else 。建议使用 SQL 字符串列表传入任务或使用分号分隔的单字符串。format_sql_for_job_facet()会在多语句模式下用分号拼接所有语句写入 OpenLineageSqlJobFacet否则仅取第一条以保持向后兼容。手动血缘标注使用inlets与outlets当算子不支持自动血缘提取或你希望覆盖自动血缘结果时可以直接在算子operator上设置inlets与outletslineage_backend_demo.py在BashOperator上设置inlets/outlets的示例展示了Dataset与Urn实体的混合使用包括envDEV、platform_instancecloud等参数以及直接放入urn:li:dataset:(...)与urn:li:dataJob:(...)URN 的写法lineage_backend_taskflow_demo.py使用 Airflow TaskFlow API 的示例。datahub_airflow_plugin.entities模块entities.py提供两类实体Dataset(platform, name, env..., platform_instanceNone)按参数构造 URN其中platform_instance可选Urn(urn:li:...)直接透传已构造好的 URN并会校验必须以urn:开头且实体类型只能是dataset或dataJobAirflow 血缘目前仅支持数据集与上游 datajob。原生 Airflow Assets/DatasetsAirflow 3.x 原生支持 Assets旧称 Datasets用于数据感知调度。DataHub 插件会自动将出现在inlets/outlets中的 Assets 捕获为血缘from airflow.sdk.definitions.asset import Asset s3_input Asset(s3://my-bucket/input/data.parquet) bigquery_output Asset(bigquery://my-project/dataset/result_table) task BashOperator( task_idprocess_data, bash_commandecho Processing, inlets[s3_input], outlets[bigquery_output], )插件按 URI scheme 将 Asset 映射到 DataHub 平台映射表定义于 _airflow_asset_adapter.py 的URI_SCHEME_TO_PLATFORMURI SchemeDataHub Platforms3://,s3a://s3gs://,gcs://gcspostgresql://postgresmysql://mysqlbigquery://bigquerysnowflake://snowflakefile://filehdfs://hdfsabfs://,abfss://adls无 scheme 的纯名称 Asset例如来自asset装饰器默认归入airflow平台。此外is_airflow_asset_alias()还能识别 Airflow 3.x 的AssetAlias有name无uri用于在执行期才确定真实 URI 的资产触发型 DAG例如 URI 中含 Jinja 模板的场景。对应的开关配置[datahub] # Set to false to disable capturing Airflow Assets as lineage (default: true) capture_airflow_assets true原生 Assets 的局限性与直接使用 DataHub 的Dataset/Urn实体相比原生 Airflow Assets 存在两点限制不支持platform_instance由 Asset URI 生成的 URN 无法包含平台实例插件只能从 URI 中提取 platform、dataset name 与 environment环境environment使用全局插件配置所有原生 Assets 都以插件配置中的cluster作为 environment无法为单个 Asset 单独指定环境。如果需要对 URN 有完全控制平台实例、每资产环境请改用 DataHub 实体类from datahub_airflow_plugin.entities import Dataset # Full control over URN components s3_input Dataset( platforms3, namemy-bucket/input/data.parquet, envPROD, platform_instanceus-west-2 # Specify platform instance ) task BashOperator( task_idprocess_data, bash_commandecho Processing, inlets[s3_input], )自定义算子场景一在自定义算子中设置 inlets/outlets如果你创建了继承自BaseOperator的自定义 Airflow 算子在重写execute时通过context[ti].task.inlets与context[ti].task.outlets设置血缘。任务运行结束后 DataHub 插件会自动拾取这些 inlets/outletsclass DbtOperator(BaseOperator): ... def execute(self, context): # do something inlets, outlets self._get_lineage() # inlets/outlets are lists of either datahub_airflow_plugin.entities.Dataset or datahub_airflow_plugin.entities.Urn context[ti].task.inlets self.inlets context[ti].task.outlets self.outlets def _get_lineage(self): # Do some processing to get inlets/outlets return inlets, outlets需要注意inlets/outlets 只能表达表级血缘如需列级血缘必须为自定义算子编写自定义 extractor。另外如果重写了pre_execute与post_execute请确保分别加上prepare_lineage与apply_lineage装饰器。仓库集成测试中给出了一个使用 SQL 解析器捕获表级血缘的自定义算子实现可直接参考custom_operator_sql_parsing.py对应 golden 文件见 tests/integration/goldens。场景二继承 SQLExecuteQueryOperator 获得自动血缘推荐这是最简单的方式——让自定义 SQL 算子继承 Airflow 的SQLExecuteQueryOperator即可自动获得✅ 内置 OpenLineage 支持✅ DataHub 增强 SQL 解析经由 SQLParser patch✅ 列级血缘提取✅ 无需编写任何血缘相关代码from typing import Any from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator class MyCustomSQLOperator(SQLExecuteQueryOperator): Custom SQL operator that inherits OpenLineage support. DataHub automatically enhances the SQL parsing with column-level lineage! def __init__(self, my_custom_param: str, **kwargs): # Add your custom parameters self.my_custom_param my_custom_param super().__init__(**kwargs) def execute(self, context: Any) - Any: # Add any custom logic before SQL execution self.log.info(fCustom param: {self.my_custom_param}) # Parent class handles SQL execution OpenLineage lineage return super().execute(context)工作原理SQLExecuteQueryOperator已经实现了get_openlineage_facets_on_complete()它内部调用 hook 的 OpenLineage 方法这些方法最终调用SQLParserDataHub 在导入时全局 patch 了SQLParser见 airflow3/_sql_parser_patch.py因此所有 SQL 解析都会被自动增强你无需编写任何血缘代码即可获得列级血缘。适用场景为标准 SQL 数据库Postgres、MySQL、Snowflake、BigQuery 等构建算子、希望集成最简单、且无需自定义血缘提取逻辑。场景三从零实现 OpenLineage 接口仅当你的算子完全无法套用SQLExecuteQueryOperator模式时才需要从零实现 OpenLineage 接口from typing import Any, Optional from airflow.models.baseoperator import BaseOperator class MyCompletelyCustomOperator(BaseOperator): For special cases where SQLExecuteQueryOperator doesnt fit. def execute(self, context: Any) - Any: # Your custom SQL execution logic pass def get_openlineage_facets_on_complete( self, task_instance: Any ) - Optional[OperatorLineage]: Implement OpenLineage interface manually. DataHubs SQLParser patch still enhances this automatically! from airflow.providers.openlineage.sqlparser import SQLParser hook self.get_db_hook() parser SQLParser( dialecthook.get_openlineage_database_dialect(hook.get_connection(self.conn_id)), default_schemahook.get_openlineage_default_schema(), ) # This uses DataHubs patched SQLParser - column lineage included! return parser.generate_openlineage_metadata_from_sql( sqlself.sql, hookhook, database_infohook.get_openlineage_database_info(hook.get_connection(self.conn_id)), )关键点DataHub 在导入时会全局 patchSQLParser.generate_openlineage_metadata_from_sql()因此任何使用 OpenLineage SQLParser 的算子都会自动获得 DataHub 增强解析含列级血缘——无论它是否继承自SQLExecuteQueryOperator。备选方案自定义算子 手动 SQL 解析如果你不希望使用 OpenLineage或运行在较旧的 Airflow 版本上可以手动用 DataHub 的 SQL 解析器提取并设置血缘from typing import Any, List, Tuple, Union from airflow.models.baseoperator import BaseOperator from datahub_airflow_plugin._config import get_enable_multi_statement from datahub_airflow_plugin._sql_parsing_common import parse_sql_with_datahub from datahub_airflow_plugin.entities import Urn class CustomSQLOperator(BaseOperator): def __init__(self, sql: Union[str, List[str]], database: str, **kwargs: Any): super().__init__(**kwargs) self.sql sql self.database database def execute(self, context: Any) - Any: # Execute SQL # ... # Extract and set lineage inlets, outlets self._get_lineage() context[ti].task.inlets inlets context[ti].task.outlets outlets def _get_lineage(self) - Tuple[List, List]: # Get multi-statement config flag enable_multi_statement get_enable_multi_statement() # Parse SQL with multi-statement support # Handles both string and list of SQL statements sql_parsing_result parse_sql_with_datahub( sqlself.sql, platformpostgres, # your platform default_databaseself.database, envPROD, default_schemaNone, graphNone, enable_multi_statementenable_multi_statement, ) inlets [Urn(table) for table in sql_parsing_result.in_tables] outlets [Urn(table) for table in sql_parsing_result.out_tables] return inlets, outletsget_enable_multi_statement()与parse_sql_with_datahub()分别封装了多语句开关的读取与 SQL 解析调用解析失败时会在日志中输出Failed to extract lineage from SQL (platform...)警告并返回空结果完整实现可查看 _sql_parsing_common.py。清理 DataHub 中过时的 Pipeline 与任务当某些 DAG 已从 Airflow 中删除、但其对应的 pipelineDataFlow与任务DataJob仍残留在 DataHub 中时这些即被称为obsolete pipelines and tasks。清理步骤如下创建一个名为Datahub_Cleanup的 DAGfrom datetime import datetime from airflow import DAG from airflow.operators.bash import BashOperator from datahub_airflow_plugin.entities import Dataset, Urn with DAG( Datahub_Cleanup, start_datedatetime(2024, 1, 1), schedule_intervalNone, catchupFalse, ) as dag: task BashOperator( task_idcleanup_obsolete_data, dagdag, bash_commandecho cleaning up the obsolete data from datahub, )摄取ingest这个 DAG。插件会依据airflow.cfg中设置的cluster值自动删除 DataHub 中该集群下所有过时的 pipeline 与任务清理 DAG 名称Datahub_Cleanup是插件内的硬编码常量见 airflow3/datahub_listener.py。查询某个 DataFlow 下的全部 DataJob如果需要查找属于某个 pipelineDataFlow的所有任务DataJob可以使用以下 GraphQL 查询query { dataFlow(urn: urn:li:dataFlow:(airflow,db_etl,prod)) { childJobs: relationships( input: { types: [IsPartOf], direction: INCOMING, start: 0, count: 100 } ) { total relationships { entity { ... on DataJob { urn } } } } } }其中urn:li:dataFlow:(airflow,db_etl,prod)的三段分别对应(platform, dag_id, cluster/env)请按你的实际集群环境替换。直接发送血缘DatahubEmitterOperator如果既无法使用插件自动提取、也不方便标注 inlets/outlets还可以在 DAG 中使用DatahubEmitterOperator直接发送血缘。完整示例见lineage_emission_dag.py核心用法如下from datahub_airflow_plugin.operators.datahub import DatahubEmitterOperator import datahub.emitter.mce_builder as builder emit_lineage_task DatahubEmitterOperator( task_idemit_lineage, datahub_conn_iddatahub_rest_default, mces[ builder.make_lineage_mce( upstream_urns[ builder.make_dataset_urn(platformsnowflake, namemydb.schema.tableA), builder.make_dataset_urn_with_platform_instance( platformsnowflake, namemydb.schema.tableB, platform_instancecloud, ), ], downstream_urnbuilder.make_dataset_urn( platformsnowflake, namemydb.schema.tableC, envDEV ), ) ], )使用该算子前必须先配置 DataHub hook。与 ingestion 类似插件同时支持 DataHub REST hook 与基于 Kafka 的 hook对应datahub-rest与datahub-kafka两类 connection可参考 hooks/datahub.py 中的 hook 实现。调试与排障缺失血缘Missing lineage如果在 DataHub 中看不到血缘请依次检查确认插件已在 Airflow 中加载进入Admin → Plugins检查 DataHub 插件是否被列出Airflow 启动时应当打印类似日志INFO [datahub_airflow_plugin.datahub_listener] DataHub plugin using DataHubRestEmitter: configured to talk to datahub_url且执行airflow plugins命令应能看到datahub_plugin且 listener 为 enabled 状态若使用插件的自动血缘确认enable_extractors配置为 true且你的算子在该功能的支持列表内若使用手动血缘标注确认 inlets/outlets 使用的是datahub_airflow_plugin.entities.Dataset或datahub_airflow_plugin.entities.Urn类。URL 生成不正确如果生成的 URL 不正确通常表现为以http://localhost:8080开头而非真实主机名需要设置 webserver 的base_url配置[webserver] base_url http://airflow.mycorp.example.comTypeError ... missing 3 required positional arguments如果出现类似下面的错误ERROR - on_task_instance_success() missing 3 required positional arguments: previous_state, task_instance, and session Traceback (most recent call last): File /home/airflow/.local/lib/python3.8/site-packages/datahub_airflow_plugin/datahub_listener.py, line 124, in wrapper f(*args, **kwargs) TypeError: on_task_instance_success() missing 3 required positional arguments: previous_state, task_instance, and session解决方案是升级acryl-datahub-airflow-plugin0.12.0.4或升级pluggy1.2.0。含.None.段的孤儿 Dataset URN如果 Airflow 任务日志中出现如下警告WARNING datahub_airflow_plugin._datahub_ol_adapter: OpenLineage Dataset name mydb.None.public.events contained None/empty segments; sanitized to mydb.public.events before producing DataHub URN. Likely upstream bug in the producer (unset field interpolated into an f-string).这说明上游 OpenLineage 产者把字面字符串None烘进了 Dataset name 的某个点号分段中。常见诱因是上游用 f-string 拼接 name 时某个插值字段为 PythonNone渲染成四个字符的None例如 Apache Airflow 的S3ToRedshiftOperator.get_openlineage_facets_on_complete使用f{database}.{self.schema}.{self.table}且没有对None做防护——任何遵循同样模式的产者Airflow 算子、dbt adapter、Spark listener、自定义集成都可能触发。DataHub 插件在生成 URN 前会自动对这些名称做 sanitize使血缘能正确对接 DataHub 原生 Redshift/Snowflake 等 ingestion 源产出的同一物理表——新血缘无需任何操作。该警告是提示性的建议针对肇事产者向上游提交 issue。几点补充行为该警告每个 worker 进程最多打印一次。首个被 sanitize 的名字会输出一条日志后续出现同名或不同名在进程生命周期内保持静默——它是警报而非流水账。若想找出全部孤儿 URN可用下面的清理命令在 DataHub 中搜索.None.或空分段。如果 Dataset name完全由None和空分段组成病态产者输出插件会保留原始名称以避免产生空 name 的 URN此时警告文本变为WARNING datahub_airflow_plugin._datahub_ol_adapter: OpenLineage Dataset name None.None had only None/empty segments; kept original to avoid emitting an empty URN. The resulting DataHub URN will literally contain None.None in its name field — search DataHub for that substring to find orphans and report the upstream producer.这类 URN 在目录中仍可检索因此可以用下面的清理命令软删除而不会消失成不可搜索的空 name URN。如果升级插件前 DataHub 中已有历史遗留的含.None.段 URN它们不会被自动清理。务必先用--dry-run预览删除结果确认无误后再真正删除# 1. Preview the URNs that would be soft-deleted (no destructive action). datahub delete --platform redshift --soft --query *.None.* --dry-run # 2. Once youve confirmed the list is correct, run the same command # without --dry-run to perform the soft-delete. datahub delete --platform redshift --soft --query *.None.*软删除的实体可以用datahub delete undo-by-filter --platform redshift恢复前提是你确实误删了超出预期的内容。调度器卡顿Scheduler stalling对于包含数千个任务的超大 Airflow 部署插件可能影响调度器性能。此时可设置环境变量DATAHUB_AIRFLOW_PLUGIN_RUN_IN_THREAD_TIMEOUT0让插件完全在后台线程运行但副作用是如果调度器在任务处理完毕后立即关闭可能丢失部分元数据。从源码看线程开关DATAHUB_AIRFLOW_PLUGIN_RUN_IN_THREAD默认即为 true超时默认 10 秒见 airflow3/datahub_listener.py。禁用 DataHub 插件提供两种禁用方式1. 通过配置禁用在airflow.cfg中设置datahub.enabledFalse然后重启 Airflow 环境以重新加载配置[datahub] enabled False2. 通过环境变量禁用Kill-Switch如果无法重启可用 kill-switch 快速禁用插件。设置环境变量AIRFLOW_VAR_DATAHUB_AIRFLOW_PLUGIN_DISABLE_LISTENERtruelistener 将立即停止处理任何事件export AIRFLOW_VAR_DATAHUB_AIRFLOW_PLUGIN_DISABLE_LISTENERtrue为何用环境变量而非 Airflow Variablelistener 钩子是在 SQLAlchemy 的after_flush事件期间主事务提交之前被调用的。此时调用Variable.get()会创建嵌套数据库会话可能干扰外层事务并造成数据丢失例如重试任务缺失 TaskInstanceHistory 记录。因此插件选择读取环境变量对应源码中的KILL_SWITCH_VARIABLE_NAME。兼容性说明DataHub 尝试在 Airflow 发布后的约 2 年内提供支持这是尽力而为的承诺——旧版本可能因依赖或安全问题无法长期维持。当前版本策略不再官方支持 Airflow 2.xApache 已对 2.x 线 EOL。要继续在 Airflow 2.x 上使用插件请固定acryl-datahub-airflow-plugin 1.6.0最后一个支持 Airflow 2 的版本Airflow 3.0 由当前版本支持。插件此前有 v1 与 v2 两代实现v2 现为默认v1 已移除。v1 的主要局限是不支持自动血缘提取。所有近期版本均要求Python 3.10。各 Airflow 版本对应的插件版本参考Airflow 1.10.xacryl-datahub-airflow-plugin 0.9.1.0v1 插件Airflow 2.0.xacryl-datahub-airflow-plugin 0.11.0.1v1 插件Airflow 2.2.xacryl-datahub-airflow-plugin 0.14.1.5v2 插件Airflow 2.3 - 2.4.3acryl-datahub-airflow-plugin 1.0.0v2 插件Airflow 2.5 与 2.6acryl-datahub-airflow-plugin 1.1.0.4v2 插件Airflow 2.7 - 2.10acryl-datahub-airflow-plugin 1.6.0v2 插件最后一个支持 Airflow 2 的版本Airflow 3.0由当前版本支持。此外DataHub 曾实现过一个 Airflow lineage backend基于 Airflow 2.2 的 lineage backend 机制其功能相当受限——不支持自动血缘提取、不捕获任务失败、且在 AWS MWAA 中无法工作因此已从代码库中移除相关文档已归档。延伸阅读Airflow 调度 DataHub ingestion用 Airflow 调度 DataHub 元数据摄取任务SQL 血缘解析器插件自动血缘所依赖的 DataHub SQL 解析机制插件源码datahub_airflow_plugin、示例 DAGexample_dags插件测试单元测试见 tests/unit集成测试 DAG 与 golden 文件见 tests/integration可用于对照验证本文所述的各种血缘行为。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表