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 模块 | 说明 |
|---|---|---|
| Operators | airflow.providers.apache.druid.operators.druid | DruidOperator,向 Druid Overlord 提交索引(ingestion)任务 |
| Hooks | airflow.providers.apache.druid.hooks.druid | DruidHook(对接 Overlord,负责提交/轮询索引任务)、DruidDbApiHook(对接 Broker,负责 SQL 查询) |
| Transfers | airflow.providers.apache.druid.transfers.hive_to_druid | HiveToDruidOperator,将 Hive 表数据导入 Druid |
| Connection types | druid(由DruidDbApiHook注册) | 在 Airflow UI 中配置 Druid 连接 |
当前发布版本为4.5.2,最低要求 Airflow>=2.11.0、pydruid>=0.6.6,支持的 Python 版本为 3.10~3.14(依据 README.rst)。
二、版本路线图:从 1.0.0 到 4.5.2 的宏观脉络
changelog.rst 记录了全部 46 个已发布版本,与 provider.yaml 中的versions列表一一对应。纵观整个演进史,可以归纳出四条主线:
- 功能线:从"仅支持提交 native batch 索引"逐步扩展到"支持 SQL 化索引(MSQ)"、SSL 校验、连接类型覆盖、SQL 查询上下文等;
- 架构线:SQL 相关类迁移至
common-sqlprovider、BaseHook迁移至 Task SDK、Hook 元数据改为 YAML 加载,与 Airflow 3.x 的组件化架构对齐; - 兼容线:最低 Airflow 版本从 2.2 一路抬升至 2.11,Python 支持从 3.7 演进到 3.14;
- 破坏性变更线:以 4.0.0 移除
DruidCheckOperator为最典型的例子。
下表按大版本归纳关键里程碑(依据 changelog 内容整理):
| 版本区间 | 核心主题 | 代表变更 |
|---|---|---|
| 1.0.0 ~ 2.x | 初始功能与打磨 | 初始版本、Check 算子重构、提交失败 Bug 修复、timeout参数、template_fields_renderers |
| 3.0.0 ~ 3.4.x | SQL 类迁移与清理 | 最低 Airflow 2.2;SQL 算子/类迁移至common-sql;移除 Python 3.7 |
| 3.5.0 ~ 3.12.x | 能力扩展 | SQL-based 索引支持、context参数、SSL 校验、日志名覆盖、连接类型覆盖、分号剥离 |
| 4.0.0 ~ 4.5.2 | Airflow 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字段,MSQ取taskId字段(见 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_time仍RUNNING,先向{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.0:
Adding optional SSL verification for druid operator (#37629),为DruidOperator增加verify_ssl参数; - 3.10.1:
Pass SSL arg to all requests in DruidOperator (#39066),修复部分请求漏传 SSL 参数的问题; - 4.3.0:
add 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.0:
Allow passing context to DruidDbApiHook (#34603),新增context参数,用于向 Druid SQL 端点传递查询上下文(如{"sqlFinalizeOuterSketches": True}),最终透传给pydruid.db.connect()的context参数; - 3.12.0:
Add possibility to override the conn type for Druid (#42793); - 4.3.0:
ssl_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.0:
Add DruidOperator template_fields_renderers fields (#19420); - 2.2.0:新增
timeout参数(由DruidHook透传); - 2.3.0:
Add 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()内将timeout、max_ingestion_time、verify_ssl、ingestion_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 连接的元数据规范:
| 参数 | 输入 |
|---|---|
| Host | Druid Broker 主机名或 IP |
| Schema | 不适用(留空) |
| Login / Password | 不适用(留空) |
| Port | Druid 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.1:
Move 'BaseHook' implementation to task SDK (#51873)、Provider Migration: Update Apache Druid for Airflow 3.0 compatibility (#52498); - 4.3.1:
Migrate 'apache/druid' provider to 'common.compat' (#57072); - 4.3.0:
Replace 'BaseHook' to Task SDK for 'apache/druid' (#52690)。
当前源码的印证是 hooks/druid.py 中的导入语句:from airflow.providers.common.compat.sdk import AirflowException, BaseHook。也就是说,DruidHook与DruidOperator不再直接依赖 Airflow 核心的BaseHook/BaseOperator,而是经由common-compat兼容层从 Task SDK 获取,这正是 changelog 4.5.2 之前一系列 "Misc" 条目背后的统一主题。
4.4 4.5.x:Hook 元数据 YAML 化与 Python 3.14
- 4.5.2:
Load hook metadata from YAML without importing Hook class (#63826)——provider 的 Hook 元数据加载不再通过导入 Hook 类完成,进一步降低启动开销并解耦依赖; - 4.5.1:
Add 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.0 | 2.2+ | — | 对应 2022 年 5 月版本 |
| 3.3.0 | 2.3+ | — | |
| 3.4.0 | 2.4+ | — | |
| 3.4.1 | — | 移除 3.7 | 明确标注 "dropped support for Python 3.7" |
| 3.6.0 | 2.5+ | — | |
| 3.7.0 | 2.6+ | — | |
| 3.10.0 | 2.7+ | — | |
| 3.11.0 | 2.8+ | — | |
| 4.0.0 | 2.9+ | — | 同时移除DruidCheckOperator |
| 4.2.0 | 2.10+ | 移除 3.9 | 同时为 pydruid 设下限 |
| 4.4.0 | 2.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.0
DruidOperator 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.1
Fix 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.1
Add support for semicolon stripping to DbApiHook, PrestoHook, and TrinoHook (#41916),允许 SQL 末尾带分号; - 依赖与打包:2.3.3
Fix 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.0、apache-airflow-providers-common-sql>=1.32.0、apache-airflow-providers-common-compat>=1.10.1、pydruid>=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_datasource、ts_dim(时间戳维度)、metric_spec(默认为[{"name": "count", "type": "count"}])、query_granularity(默认NONE,即毫秒级)、segment_granularity(默认DAY)等,均在construct_ingest_query()中组装为完整的 Druid 索引规格(见 hive_to_druid.py)。
八、升级与选型建议(基于 changelog 的实践结论)
- 若你仍在使用
DruidCheckOperator:请升级到 4.0.0+ 之前先迁移至SQLCheckOperator,这是 4.0.0 明确的破坏性变更; - 若你使用 Druid SQL 查询:统一走
SQLExecuteQueryOperator+druid连接类型,Druid 专用 SQL 算子已在 3.1.0/3.2.0 迁入 common-sql,不再需要单独维护; - 关注最低版本门槛:当前 4.5.2 要求 Airflow 2.11+,升级 Airflow 时需同步评估 provider 版本边界(见第五节表格);
- 生产环境务必设置
max_ingestion_time:源码层面提供了"超时自动 shutdown"的保护,避免索引任务无限挂起占用资源; - HTTPS 场景配置好 SSL 链路:
verify_ssl、ssl_verify_cert、ca_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),仅供参考