Apache Airflow 中使用 SQLExecuteQueryOperator 连接与查询 Apache Druid 完整指南

发布时间:2026/9/13 18:49:19
Apache Airflow 中使用 SQLExecuteQueryOperator 连接与查询 Apache Druid 完整指南 Apache Airflow 中使用 SQLExecuteQueryOperator 连接与查询 Apache Druid 完整指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读本文讲解如何在 Apache Airflow 中通过通用的SQLExecuteQueryOperator对 Apache Druid 集群执行 SQL 查询覆盖从 provider 安装、Airflow 连接配置、示例 DAG 编写到底层DruidDbApiHook工作原理的完整链路。读者完成后将能够正确配置指向 Druid broker 的 Airflow Connection、使用一条conn_id运行任意 Druid SQL包括sys.segments元数据查询与INFORMATION_SCHEMA表结构查询并理解参数优先级规则从而替代已弃用的专用 Druid 查询算子。为什么用 SQLExecuteQueryOperator 而非专用 Druid 算子在 Apache Airflow 的 Apache Druid providerapache-airflow-providers-apache-druid中执行 SQL 查询的官方推荐方式是复用 common SQL provider 提供的通用算子SQLExecuteQueryOperator完整类路径为airflow.providers.common.sql.operators.sql.SQLExecuteQueryOperator而不是维护一个 Druid 专属的查询算子。依据 operators.rst 中的说明以前可能使用过一个专门为 Druid 设计的算子在弃用之后请改用SQLExecuteQueryOperator。这意味着查询 Druid broker 走的是 Airflow 统一的 DB-API 抽象层与查询 MySQL、Postgres 等数据库共用同一套算子语义sql、autocommit、parameters、handler、return_last等参数该 provider 仍保留DruidOperator与DruidHook但它们只服务于索引ingestion任务提交不再承担 SQL 查询职责两者职责边界清晰查询用SQLExecuteQueryOperator提交批式/ MSQ 索引任务用DruidOperator见 operators/druid.py。安装 provider 与版本要求文档明确要求先安装 Druid provider 包才能启用 Druid 支持pip install apache-airflow-providers-apache-druid根据 index.rstRelease 4.5.2记录的核心依赖依赖包版本要求apache-airflow2.11.0apache-airflow-providers-common-sql1.32.0提供SQLExecuteQueryOperator与DbApiHookapache-airflow-providers-common-compat1.10.1pydruid0.6.6Druid 的 DB-API 驱动DruidDbApiHook.get_conn()正是通过pydruid.db.connect(...)建立连接其中pydruid是底层 SQL 驱动负责把标准 DB-API 调用翻译成 Druid broker 的/druid/v2/sql/HTTP 查询common-sqlprovider 则提供算子与 Hook 基类。如需同时使用 Hive 到 Druid 的数据传输功能可额外安装pip install apache-airflow-providers-apache-druid[apache.hive]配置指向 Druid broker 的 Airflow Connection使用conn_id参数连接 Druid 实例时连接元数据需按如下结构填写见 operators.rst 的连接元数据表参数填写内容Host: stringDruid broker 的主机名或 IP 地址Schema: string不适用留空Login: string不适用留空Password: string不适用留空Port: intDruid broker 端口默认8082Extra: JSON额外的连接配置例如{endpoint: /druid/v2/sql/, method: POST, ssl_verify_cert: false}关键点说明查询目标是 Druid broker默认端口8082而不是 Overlord8081用于索引任务提交Extra中最重要的是endpoint它指定 Druid SQL 查询的 HTTP 路径。从 hooks/druid.py 的DruidDbApiHook.get_conn()实现可见其默认值为/druid/v2/sql不带尾斜杠pathconn.extra_dejson.get(endpoint, /druid/v2/sql), schemeconn.extra_dejson.get(schema, http), userconn.login, passwordconn.password, ssl_verify_certconn.extra_dejson.get(ssl_verify_cert, True),同一 Extra 中还支持schema协议默认http可设为https与ssl_verify_cert是否校验 TLS 证书默认True以及可选contextDruid SQL 查询上下文参数例如{sqlFinalizeOuterSketches: True}见DruidDbApiHook.__init__的context参数说明连接类型应选择druidDruidDbApiHook中定义了conn_type druid此时get_uri()会生成形如druid://localhost:8082/druid/v2/sql/的连接串。编写示例 DAG查询 datasource 元数据官方系统测试 DAG example_druid.py 给出了三个最典型的 Druid SQL 查询场景完整示例即文档中[START howto_operator_druid]至[END howto_operator_druid]之间的内容如下from __future__ import annotations import datetime from textwrap import dedent from airflow import DAG from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator DAG_ID example_druid with DAG( dag_idDAG_ID, start_datedatetime.datetime(2025, 1, 1), default_args{conn_id: my_druid_conn}, scheduleonce, catchupFalse, ) as dag: # Task: List all published datasources in Druid. list_datasources_task SQLExecuteQueryOperator( task_idlist_datasources, sqlSELECT DISTINCT datasource FROM sys.segments WHERE is_published 1, ) # Task: Describe the schema for the wikipedia datasource. # Note: This query returns column information if the datasource exists. describe_wikipedia_task SQLExecuteQueryOperator( task_iddescribe_wikipedia, sqldedent( SELECT COLUMN_NAME, DATA_TYPE FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME wikipedia ).strip(), ) # Task: Count rows for the wikipedia datasource. # Here we count the segments for wikipedia. If the datasource is not ingested, it returns 0. select_count_from_datasource SQLExecuteQueryOperator( task_idselect_count_from_datasource, sqlSELECT COUNT(*) FROM sys.segments WHERE datasource wikipedia, ) list_datasources_task describe_wikipedia_task select_count_from_datasource要点解析conn_id既可以直接传给算子也可以像示例一样放在default_args中二者等价算子显式传入时优先三个任务覆盖三类常见需求枚举已发布 datasourcesys.segments元数据表、查看 datasource 的列与类型INFORMATION_SCHEMA.COLUMNS、统计 datasource 行数/segment 数对未摄入的 datasource 返回 0从 sql.py 中SQLExecuteQueryOperator的源码可知sql与parameters均支持模板渲染template_fields (sql, parameters, ...)template_ext (.sql, .json)因此 SQL 既可以内联字符串也可以指向以.sql结尾的模板文件其余常用参数包括autocommit默认False、handler默认fetch_all_handler即拉取全部结果、return_last默认True只返回最后一条语句的结果、split_statements、show_return_value_in_logs默认False谨慎打印大结果集以及do_xcom_push开启后查询结果会推送到 XCom。参数优先级规则operators.rst 明确给出了参数优先级约定直接通过SQLExecuteQueryOperator()提供的参数优先于 Airflow 连接元数据中指定的参数如schema、login、password等。即当算子的显式参数与 Connection 中的元数据冲突时以算子参数为准。这一约定与DbApiHook的通用行为一致——算子层参数如database会覆盖连接定义是 Airflow 中“代码优先于配置”的体现。实际使用中建议将 broker 地址等基础设施信息放在 Connection 里而把每次查询差异化的内容SQL、参数、return_last、handler等放在任务定义中。底层原理DruidDbApiHook 如何工作SQLExecuteQueryOperator本身不感知 Druid它通过连接类型自动匹配到DruidDbApiHook定义于 hooks/druid.py。该 Hook 继承自 common SQL 的DbApiHook关键实现如下连接建立get_conn()调用pydruid.db.connect(host..., port..., path..., scheme..., user..., password..., context..., ssl_verify_cert...)将 Airflow Connection 的字段逐一映射为 pydruid 驱动参数Extra 中的endpoint→path、schema→scheme、ssl_verify_cert→ssl_verify_cert查询上下文context参数透传给 pydruid 的connect()用于向 Druid SQL 端点传递查询上下文如 sketch 合并策略、查询超时等标准 DB-API 能力DruidDbApiHook实现了get_first、get_records、get_df支持pandas与polars两种df_type、get_df_by_chunks等方法因此SQLExecuteQueryOperator的全部输出处理逻辑handler、output_processor、XCom 推送都可以直接复用不支持的操作该 Hook 明确抛出NotImplementedError的方法包括set_autocommit与insert_rowssupports_autocommit False意味着对 Druid 不应依赖事务自动提交语义也不应尝试用该 Hook 做行级写入——写入应走 Druid 的摄入ingestion流程。单元测试 test_druid.py 中的TestDruidDbApiHook验证了上述映射例如在 Extra 为{endpoint: /test/endpoint, schema: https}时mock_connect收到path/test/endpoint、schemehttps在 Extra 含ssl_verify_cert: False时驱动收到ssl_verify_certFalse证明连接层配置的完整传递链路。与索引算子的职责划分了解即可Druid provider 还保留了面向数据摄入的组件与本文的 SQL 查询算子职责互补避免混淆DruidHookhooks/druid.py面向 Druid Overlord 提交索引任务支持IngestionType.BATCHNative batchendpoint 取自 Extra 的endpoint与IngestionType.MSQSQL-based ingestionendpoint 取自 Extra 的msq_endpoint并轮询任务状态RUNNING/SUCCESS/FAILED超过max_ingestion_time会主动发起 shutdownDruidOperatoroperators/druid.py封装DruidHook接收json_index_file支持.json模板渲染示例见 example_druid_dag.pyHiveToDruidOperatortransfers/hive_to_druid.py将 Hive 表数据转存到 Druid。因此在你的 DAG 中读 Druid 用SQLExecuteQueryOperator本文主题写 Druid 用DruidOperator/HiveToDruidOperator二者分别对应 broker 查询端点默认 8082与 Overlord 索引端点默认 8081连接配置也应按用途分别创建。参考文档Druid provider 使用指南与操作符说明operators.rstDruid provider 包总览、依赖与安装说明index.rst官方系统测试示例 DAGexample_druid.py 与 example_druid_dag.py底层 Hook 实现hooks/druid.py索引算子实现operators/druid.py通用 SQL 算子源码sql.pyHook 单元测试test_druid.py【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询