Apache Airflow Apache Druid Provider 演进全解析:从 changelog 到源码的版本路线图与迁移指南
2026/9/13 9:12:42 网站建设 项目流程

Apache Airflow Apache Druid Provider 演进全解析:从 changelog 到源码的版本路线图与迁移指南

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

本篇技术指南以当前仓库中 Apache Druid Provider 官方变更日志 为骨架,结合 provider 元数据、Hook/Operator 源码 与系统测试 DAG,系统梳理apache-airflow-providers-apache-druid从 1.0.0 到 4.5.2 的功能演进、架构重构与破坏性变更。读完你将掌握该 provider 的能力边界、核心实现原理、与 Airflow 核心/SDK 的兼容关系,以及从旧版本升级时的关键注意事项。

一、Provider 包概览:它是什么,能做什么

apache-airflow-providers-apache-druid是 Apache Airflow 官方社区维护的 Provider 包,用于将 Apache Druid(实时分析型数据库)接入 Airflow 工作流。根据 get_provider_info.py 与 provider.yaml 的声明,该包当前提供四类能力:

能力类别提供的 Python 模块说明
Operatorsairflow.providers.apache.druid.operators.druidDruidOperator,向 Druid Overlord 提交索引(ingestion)任务
Hooksairflow.providers.apache.druid.hooks.druidDruidHook(对接 Overlord,负责提交/轮询索引任务)、DruidDbApiHook(对接 Broker,负责 SQL 查询)
Transfersairflow.providers.apache.druid.transfers.hive_to_druidHiveToDruidOperator,将 Hive 表数据导入 Druid
Connection typesdruid(由DruidDbApiHook注册)在 Airflow UI 中配置 Druid 连接

当前发布版本为4.5.2,最低要求 Airflow>=2.11.0pydruid>=0.6.6,支持的 Python 版本为 3.10~3.14(依据 README.rst)。

二、版本路线图:从 1.0.0 到 4.5.2 的宏观脉络

changelog.rst 记录了全部 46 个已发布版本,与 provider.yaml 中的versions列表一一对应。纵观整个演进史,可以归纳出四条主线:

  1. 功能线:从"仅支持提交 native batch 索引"逐步扩展到"支持 SQL 化索引(MSQ)"、SSL 校验、连接类型覆盖、SQL 查询上下文等;
  2. 架构线:SQL 相关类迁移至common-sqlprovider、BaseHook迁移至 Task SDK、Hook 元数据改为 YAML 加载,与 Airflow 3.x 的组件化架构对齐;
  3. 兼容线:最低 Airflow 版本从 2.2 一路抬升至 2.11,Python 支持从 3.7 演进到 3.14;
  4. 破坏性变更线:以 4.0.0 移除DruidCheckOperator为最典型的例子。

下表按大版本归纳关键里程碑(依据 changelog 内容整理):

版本区间核心主题代表变更
1.0.0 ~ 2.x初始功能与打磨初始版本、Check 算子重构、提交失败 Bug 修复、timeout参数、template_fields_renderers
3.0.0 ~ 3.4.xSQL 类迁移与清理最低 Airflow 2.2;SQL 算子/类迁移至common-sql;移除 Python 3.7
3.5.0 ~ 3.12.x能力扩展SQL-based 索引支持、context参数、SSL 校验、日志名覆盖、连接类型覆盖、分号剥离
4.0.0 ~ 4.5.2Airflow 3.x 对齐移除DruidCheckOperator;Task SDK 迁移;最低 Airflow 2.11;YAML Hook 元数据加载

三、核心功能演进深度解析(结合源码)

3.1 索引提交:从 BATCH 到 MSQ(3.5.0 引入 SQL-based 任务支持)

changelog 3.5.0 记录DruidHook add SQL-based task support (#32795),这是 provider 能力的一次重要扩展。在 hooks/druid.py 中可以看到对应的IngestionType枚举:

class IngestionType(Enum): """Druid Ingestion Type. Could be Native batch ingestion or SQL-based ingestion.""" BATCH = 1 MSQ = 2

两种类型在源码中的差异非常清晰:

  • URL 构造不同get_conn_url()中,BATCH使用连接 Extra 里的endpoint(默认空,即指向 Overlord 根路径),而MSQ使用msq_endpoint(见 hooks/druid.py);
  • 任务 ID 字段不同submit_indexing_job()中,BATCH从响应取task字段,MSQtaskId字段(见 hooks/druid.py)。

3.2 任务状态查询 URL(4.1.0 新增)

changelog 4.1.0 记录Add method to retrieve Druid task status URL based on ingestion type (#47238)。对应源码为get_status_url()(见 hooks/druid.py):对MSQ类型,会优先使用连接 Extra 中status_endpoint(默认druid/indexer/v1/task),并允许通过 Extra 的schema覆盖协议;对BATCH类型则复用get_conn_url()。这一方法被submit_indexing_job()用来拼出形如.../druid/indexer/v1/task/{taskId}/status的轮询地址。

3.3 轮询与超时控制机制

submit_indexing_job()(见 hooks/druid.py)完整实现了"提交 → 轮询 → 判定结果"的闭环:

  • 提交:向 Overlord 发POST,仅接受200 <= code < 300的响应,否则抛出AirflowException
  • 轮询:按timeout秒(构造时强制>= 1,否则抛ValueError)间隔查询任务状态,日志输出已运行秒数;
  • 超时保护:若超过max_ingestion_timeRUNNING,先向{url}/{task_id}/shutdown发起关闭请求,再抛出异常——对应 changelog 3.8.0 的 "Fix successful Apache Druid task submissions reported as failed (#36813)" 以及 2.2.0 的Add timeout parameter to DruidOperator (#19984)
  • 状态判定:SUCCESS结束轮询,FAILED抛异常,其他状态视为无法识别。

3.4 SSL 验证能力(3.9.0 引入,3.10.1/4.3.0 完善)

SSL 相关的演进在 changelog 中出现了三次:

  • 3.9.0Adding optional SSL verification for druid operator (#37629),为DruidOperator增加verify_ssl参数;
  • 3.10.1Pass SSL arg to all requests in DruidOperator (#39066),修复部分请求漏传 SSL 参数的问题;
  • 4.3.0add ssl_verify_cert support to DruidDbApiHook.get_conn (#52926),将校验能力延伸到 SQL 查询链路。

当前源码中的落点是两处:

  • DruidHook.get_verify()(见 hooks/druid.py):当verify_ssl=False且连接 Extra 配置了ca_bundle_path时,返回该 CA 包路径用于校验;否则原样返回verify_ssl
  • DruidDbApiHook.get_conn()(见 hooks/druid.py):通过conn.extra_dejson.get("ssl_verify_cert", True)传入pydruid.db.connect()

此外,DruidHook.get_auth()(见 hooks/druid.py)支持从连接配置读取login/password构造 HTTP Basic Auth,用于对接启用了druid-basic-security扩展的 Druid 集群——这一点在 changelog 2.0.1Fix error in Druid connection attribute retrieval (#17095)中也有体现。

3.5 SQL 查询:DruidDbApiHook 的能力集合(3.6.0 / 4.3.0)

DruidDbApiHook继承自DbApiHook(common-sql),专门用于查询 Druid Broker。其关键能力演进:

  • 3.6.0Allow passing context to DruidDbApiHook (#34603),新增context参数,用于向 Druid SQL 端点传递查询上下文(如{"sqlFinalizeOuterSketches": True}),最终透传给pydruid.db.connect()context参数;
  • 3.12.0Add possibility to override the conn type for Druid (#42793)
  • 4.3.0ssl_verify_cert支持(见上文)。

get_conn()从连接 Extra 中读取endpoint(默认/druid/v2/sql)与schema(默认http),从而构造出 pydruid 连接;get_uri()则输出形如druid://localhost:8082/druid/v2/sql/的 URI。

3.6 DruidOperator 的模板化能力(2.1.0 / 2.3.0 / 2.2.0)

DruidOperator(见 operators/druid.py)负责读取 JSON 索引规格并提交给DruidHook。changelog 中与模板化相关的演进:

  • 2.1.0Add DruidOperator template_fields_renderers fields (#19420)
  • 2.2.0:新增timeout参数(由DruidHook透传);
  • 2.3.0Add more SQL template fields renderers (#21237)

当前源码中的模板配置为:

template_fields: Sequence[str] = ("json_index_file",) template_ext: Sequence[str] = (".json",) template_fields_renderers = {"json_index_file": "json"}

json_index_file支持 Jinja 模板渲染,且支持以.json文件形式提供(Airflow 会自动读取文件内容并渲染)。execute()内将timeoutmax_ingestion_timeverify_sslingestion_type全部透传给DruidHook后调用submit_indexing_job()

四、架构性重构与破坏性变更(升级必读)

4.1 4.0.0:移除 DruidCheckOperator(破坏性变更)

changelog 4.0.0 是唯一明确标注Breaking changes的近期版本,核心内容是:

所有已废弃的类、参数与特性已从 Apache Druid provider 中移除。DruidCheckOperator已被移除,请改用airflow.providers.common.sql.operators.sql.SQLCheckOperator

这正是 provider 生态"SQL 类统一收敛到 common-sql"战略的延续——早在 3.1.0(Move all SQL classes to common-sql provider)与 3.2.0(Move all "old" SQL operators to common.sql providers)就已开始铺垫。升级到 4.0.0 及以上的用户必须将 DAG 中的DruidCheckOperator替换为SQLCheckOperator,否则会因类不存在而导入失败。

4.2 SQL 查询算子迁移到 SQLExecuteQueryOperator

operators.rst 明确指出:对 Druid 集群执行 SQL 查询应使用airflow.providers.common.sql.operators.sql.SQLExecuteQueryOperator,而非专用的 Druid 算子。文档还给出了 Druid 连接的元数据规范:

参数输入
HostDruid Broker 主机名或 IP
Schema不适用(留空)
Login / Password不适用(留空)
PortDruid Broker 端口(默认 8082)
Extra (JSON){"endpoint": "/druid/v2/sql/", "method": "POST", "ssl_verify_cert": false}

同时强调:直接传给SQLExecuteQueryOperator的参数优先于连接元数据中的同名配置。完整可运行的示例见 example_druid.py,其中展示了查询已发布数据源、查询列信息、统计 segment 数三个典型任务:

list_datasources_task = SQLExecuteQueryOperator( task_id="list_datasources", sql="SELECT DISTINCT datasource FROM sys.segments WHERE is_published = 1", ) describe_wikipedia_task = SQLExecuteQueryOperator( task_id="describe_wikipedia", sql=dedent(""" SELECT COLUMN_NAME, DATA_TYPE FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = 'wikipedia' """).strip(), ) select_count_from_datasource = SQLExecuteQueryOperator( task_id="select_count_from_datasource", sql="SELECT COUNT(*) FROM sys.segments WHERE datasource = 'wikipedia'", )

4.3 BaseHook 向 Task SDK 迁移(4.2.1 / 4.3.1)

面向 Airflow 3.x 的组件化重构在 changelog 中清晰可见:

  • 4.2.1Move 'BaseHook' implementation to task SDK (#51873)Provider Migration: Update Apache Druid for Airflow 3.0 compatibility (#52498)
  • 4.3.1Migrate 'apache/druid' provider to 'common.compat' (#57072)
  • 4.3.0Replace 'BaseHook' to Task SDK for 'apache/druid' (#52690)

当前源码的印证是 hooks/druid.py 中的导入语句:from airflow.providers.common.compat.sdk import AirflowException, BaseHook。也就是说,DruidHookDruidOperator不再直接依赖 Airflow 核心的BaseHook/BaseOperator,而是经由common-compat兼容层从 Task SDK 获取,这正是 changelog 4.5.2 之前一系列 "Misc" 条目背后的统一主题。

4.4 4.5.x:Hook 元数据 YAML 化与 Python 3.14

  • 4.5.2Load hook metadata from YAML without importing Hook class (#63826)——provider 的 Hook 元数据加载不再通过导入 Hook 类完成,进一步降低启动开销并解耦依赖;
  • 4.5.1Add Python 3.14 Support (#63520),与 README.rst 中列出的 3.10~3.14 支持范围一致。

五、支持矩阵演进:Airflow 与 Python 版本边界

changelog 中多次出现 "This release of provider is only available for Airflow X+" 的说明,整理成时间线如下(依据 changelog 各版本 note):

版本最低 Airflow最低 Python备注
3.0.02.2+对应 2022 年 5 月版本
3.3.02.3+
3.4.02.4+
3.4.1移除 3.7明确标注 "dropped support for Python 3.7"
3.6.02.5+
3.7.02.6+
3.10.02.7+
3.11.02.8+
4.0.02.9+同时移除DruidCheckOperator
4.2.02.10+移除 3.9同时为 pydruid 设下限
4.4.02.11+当前 4.5.2 沿用该边界

另有若干 Python 版本相关条目散见于各版本:2.3.1Support for Python 3.10、4.2.1Drop support for Python 3.9 (#52072)、4.3.0Add Python 3.13 support。需要特别提醒:Provider 的最低 Airflow 版本由 Apache Airflow 社区维护的 PROVIDERS.rst 中的支持策略统一约束,从 3.0.0 到 4.4.0 每次抬升都对应社区对旧 Airflow 版本停止支持的时间点。

六、值得关注的 Bug 修复与质量改进

changelog 中记录了大量"看不见但很重要"的修复,按主题归纳:

  • 提交与状态判定:1.1.0DruidOperator fails to submit ingestion tasks (#14418);3.8.0Fix successful Apache Druid task submissions reported as failed (#36813);3.3.0BugFix - Druid Airflow Exception to about content (#27174)
  • 连接属性:2.0.1Fix error in Druid connection attribute retrieval (#17095);3.10.2Clean up remaining getattr connection DbApiHook (#40665)
  • SSL 传递:3.9.0 与 3.10.1 的两次修复(见 3.4 节);
  • SQL 行为:3.12.1Add support for semicolon stripping to DbApiHook, PrestoHook, and TrinoHook (#41916),允许 SQL 末尾带分号;
  • 依赖与打包:2.3.3Fix mistakenly added install_requires for all providers (#22382);3.3.1Bump common.sql provider to 1.3.1;4.2.0Lower bind pyspark and pydruid to relatively new versions (#50205)
  • 工程质量:3.8.1 将所有类/函数/方法的弃用声明切换为装饰器;4.1.1 移除冗余else块并改进示例文档。

七、从 changelog 到实战:当前版本的完整工作流

7.1 安装与依赖

依据 README.rst:

pip install apache-airflow-providers-apache-druid # 如需使用 Hive→Druid 传输算子,额外安装 cross-provider 依赖: pip install apache-airflow-providers-apache-druid[apache.hive]

硬性依赖为apache-airflow>=2.11.0apache-airflow-providers-common-sql>=1.32.0apache-airflow-providers-common-compat>=1.10.1pydruid>=0.6.6

7.2 两种典型任务形态

形态一:提交索引任务(DruidOperator)

DruidOperator( task_id="submit_ingestion", json_index_file="/path/to/index_spec.json", # 支持 Jinja 模板与 .json 文件渲染 druid_ingest_conn_id="druid_ingest_default", timeout=1, # 轮询间隔,必须 >= 1 max_ingestion_time=3600, # 超时后自动 shutdown 任务并失败 ingestion_type=IngestionType.BATCH, # 或 IngestionType.MSQ verify_ssl=True, # False 时可用 ca_bundle_path 指定 CA )

形态二:SQL 查询(SQLExecuteQueryOperator + DruidDbApiHook)

配置好druid类型连接(Broker 端口默认 8082)后,直接使用SQLExecuteQueryOperator,示例见本文 4.2 节;如需传递 Druid SQL 查询上下文,可在 Hook 层通过context参数注入。

7.3 Hive → Druid 数据迁移

HiveToDruidOperator(见 transfers/hive_to_druid.py)是 changelog 之外该 provider 保留的另一项核心能力,其执行流程为:用HiveCliHook将 SQL 查询结果落为 HDFS 上的 TSV 临时表 → 用HiveMetastoreHook读取列结构与 HDFS 路径 → 用DruidHook提交index_hadoop类型的原生批索引任务 →finally中清理临时表。关键参数包括druid_datasourcets_dim(时间戳维度)、metric_spec(默认为[{"name": "count", "type": "count"}])、query_granularity(默认NONE,即毫秒级)、segment_granularity(默认DAY)等,均在construct_ingest_query()中组装为完整的 Druid 索引规格(见 hive_to_druid.py)。

八、升级与选型建议(基于 changelog 的实践结论)

  1. 若你仍在使用DruidCheckOperator:请升级到 4.0.0+ 之前先迁移至SQLCheckOperator,这是 4.0.0 明确的破坏性变更;
  2. 若你使用 Druid SQL 查询:统一走SQLExecuteQueryOperator+druid连接类型,Druid 专用 SQL 算子已在 3.1.0/3.2.0 迁入 common-sql,不再需要单独维护;
  3. 关注最低版本门槛:当前 4.5.2 要求 Airflow 2.11+,升级 Airflow 时需同步评估 provider 版本边界(见第五节表格);
  4. 生产环境务必设置max_ingestion_time:源码层面提供了"超时自动 shutdown"的保护,避免索引任务无限挂起占用资源;
  5. HTTPS 场景配置好 SSL 链路verify_sslssl_verify_certca_bundle_path三个参数分别覆盖索引提交与 SQL 查询两条链路,缺一不可。

九、延伸阅读

  • 变更日志原文:providers/apache/druid/docs/changelog.rst
  • 详细提交列表:providers/apache/druid/docs/commits.rst
  • Operator 使用指南:providers/apache/druid/docs/operators.rst
  • 安全说明:providers/apache/druid/docs/security.rst
  • 核心源码:hooks/druid.py、operators/druid.py、transfers/hive_to_druid.py
  • Provider 元数据与依赖声明:provider.yaml、README.rst
  • 单元测试:hooks/test_druid.py、operators/test_druid.py

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询