news 2026/9/13 15:49:23

Apache Airflow 集成 Apache Pinot 实战:使用 SQLExecuteQueryOperator 执行实时 OLAP 查询

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow 集成 Apache Pinot 实战:使用 SQLExecuteQueryOperator 执行实时 OLAP 查询

Apache Airflow 集成 Apache Pinot 实战:使用 SQLExecuteQueryOperator 执行实时 OLAP 查询

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

Apache Airflow 官方并没有为 Apache Pinot 提供专门的 Operator,而是统一推荐使用通用 SQL 执行算子SQLExecuteQueryOperator来对 Pinot Broker 发起标准 SQL 查询。本文以仓库中的 operators.rst 为核心,完整讲解连接配置、示例 DAG、底层 Hook 调用链与参数优先级,帮助你快速在 Airflow 中构建针对 Pinot 的查询任务,并掌握可复用的最佳实践。

为什么 Pinot 没有专属 Operator

Apache Pinot 是一个面向列的、分布式的开源 OLAP 数据存储(由 Java 编写),专为低延迟分析场景设计,适合在不可变数据上进行聚合类快速分析,也支持实时数据摄入。Pinot 对外暴露的是标准 SQL 查询接口,因此 Airflow 社区将其接入方式统一收敛到通用 SQL 体系:

  • 使用SQLExecuteQueryOperator(位于airflow.providers.common.sql.operators.sql)执行查询;
  • 底层由 Pinot Provider 提供的PinotDbApiHook负责建立与 Pinot Broker 的连接。

这一点在文档中有两条明确的note说明:

  1. Apache Pinot 没有专属 Operator,请直接使用SQLExecuteQueryOperator
  2. 必须先安装对应的 Provider 包(apache-airflow-providers-apache-pinot),才能启用 Apache Pinot 支持。

从源码角度看,SQLExecuteQueryOperator继承自BaseSQLOperator,在execute阶段通过get_db_hook()获取对应的 DB Hook,再调用hook.run(...)执行 SQL(见 sql.py)。当conn_id指向一个pinot类型的连接时,框架会自动解析到PinotDbApiHook,这就是"无专属 Operator 却能开箱即用"的实现原理。

前置条件:安装 Provider 包

在已有的 Airflow 环境中执行:

pip install apache-airflow-providers-apache-pinot

根据 index.rst 中的要求,该 Provider 的最低依赖如下:

PIP 包版本要求
apache-airflow>=2.11.0
apache-airflow-providers-common-compat>=1.10.1
apache-airflow-providers-common-sql>=1.32.0
pinotdb>=5.1.0

其中pinotdb是 Pinot 官方的 Python DB-API 驱动,PinotDbApiHook正是通过它连接 Broker。

配置 Pinot 连接(Connection)

使用SQLExecuteQueryOperator时,通过conn_id参数指向一个名为pinot类型的 Airflow Connection。连接元数据的结构如下:

参数取值
Host: stringPinot Broker 的主机名或 IP 地址
Port: intPinot Broker 端口(默认:8000)
Schema: string不使用
Extra: JSON可选字段,例如{"endpoint": "query/sql"}

注意事项:

  • 文档表格中默认端口写为 8000,但实际生产环境常见的 Broker 端口是 8099 或 9000 等,请以你实际部署的 Pinot Broker 端口为准
  • Schema字段在查询场景下不被使用,但PinotDbApiHook.get_conn()会读取extra_dejson中的schema键作为请求协议(scheme),默认http
  • Extra中的endpoint用于指定 Broker 的 SQL 查询端点,默认值为/query/sql

更深入的实现细节可以参考 pinot.py 中PinotDbApiHook.get_conn()的源码:它通过pinotdb.connect(host=..., port=..., username=conn.login, password=conn.password, path=..., scheme=...)建立连接。也就是说,Airflow Connection 的LoginPassword字段会被透传给 Pinot 做身份认证,Extra中支持的键包括:

  • endpoint:SQL 查询端点路径,默认query/sql
  • schema:连接协议,默认http(对应conn_type)。

此外,PinotDbApiHook.get_uri()会拼出形如http://localhost:9000/query/sql的 URI(conn_type://[login:password@]host:port/endpoint),可用于日志展示或诊断。

使用 SQLExecuteQueryOperator 编写查询 DAG

文档通过exampleinclude指令嵌入了系统测试示例 example_pinot.py,下面是其核心内容(对应[START howto_operator_pinot][END howto_operator_pinot]之间的代码):

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_pinot" with DAG( dag_id=DAG_ID, start_date=datetime.datetime(2025, 1, 1), default_args={"conn_id": "my_pinot_conn"}, schedule="@once", catchup=False, ) as dag: # Task: Simple query to test connection and query engine select_1_task = SQLExecuteQueryOperator( task_id="select_1", sql="SELECT 1", ) # Task: Count total records in airlineStats (sample table) count_airline_stats = SQLExecuteQueryOperator( task_id="count_airline_stats", sql="SELECT COUNT(*) FROM airlineStats", ) # Task: Group by Carrier and count flights group_by_carrier = SQLExecuteQueryOperator( task_id="group_by_carrier", sql=dedent(""" SELECT Carrier, COUNT(*) AS flight_count FROM airlineStats GROUP BY Carrier ORDER BY flight_count DESC LIMIT 5 """).strip(), ) select_1_task >> count_airline_stats >> group_by_carrier

代码要点拆解

  1. conn_id 的两种指定方式:示例在 DAG 的default_args中统一设置了"conn_id": "my_pinot_conn",因此每个SQLExecuteQueryOperator无需重复传参;你也可以在单个 Task 上通过conn_id=显式覆盖。
  2. sql参数支持多行 SQLSQLExecuteQueryOperatortemplate_fields = ("sql", "parameters", ...),说明sql是支持模板渲染的字段。多行语句可以用textwrap.dedent或三引号书写,保持可读性。
  3. 任务编排:示例用位运算>>将三个查询串成线性链:select_1(连通性测试)→count_airline_stats(全表计数)→group_by_carrier(分组聚合)。其中SELECT 1常被用来快速验证连接与查询引擎是否正常。

算子关键参数速查

结合SQLExecuteQueryOperator的源码定义(sql.py),除conn_id外,最常用的参数有:

参数默认值说明
sql必填要执行的 SQL 字符串或指向.sql/.json模板文件的路径
autocommitFalse是否自动提交(Pinot 为只读查询场景,通常无需关心)
parametersNone用于渲染 SQL 的参数字典/序列
handlerfetch_all_handler应用到 cursor 的结果处理函数
split_statementsNone是否按语句拆分执行,默认沿用 Hook 的run方法行为
return_lastTrue多条语句时仅返回最后一条的结果
show_return_value_in_logsFalse是否将算子输出打印到任务日志(谨慎用于大数据集)
requires_result_fetchFalse是否强制在完成前抓取查询结果
do_xcom_push继承自 BaseOperatorTrue时结果会自动写入 XCom,供下游任务消费

注意:PinotDbApiHook声明了supports_autocommit = False,且其set_autocommitinsert_rows均直接抛出NotImplementedError——这印证了 Pinot 集成定位是只读分析查询,不适合通过该 Hook 写入数据。

参数优先级:算子入参 > 连接元数据

文档末尾特别强调:

Parameters provided directly viaSQLExecuteQueryOperator()take precedence over those specified in the Airflow connection metadata.

直接在SQLExecuteQueryOperator()中传入的参数,优先级高于 Airflow 连接元数据中的配置。这一点是 Airflow 通用 SQL 体系的统一行为:Connection 提供的是"默认连接信息",而每个 Task 上的显式参数可以在不修改 Connection 的前提下覆盖默认值。例如可以在 Connection 中配置默认的 Broker 地址,而在特定 Task 上通过参数临时指向其他 Broker 实例。

底层链路:从 Operator 到 Pinot Broker

一次查询任务的完整调用链可以概括为:

SQLExecuteQueryOperator.execute() │ 1. get_db_hook() 解析 conn_id(pinot 类型 → PinotDbApiHook) ▼ PinotDbApiHook.run(sql, ...) # 继承自 DbApiHook │ 2. get_conn() 调用 pinotdb.connect(...) ▼ pinotdb.Connection → Cursor.execute(sql) │ 3. 向 http://<host>:<port>/query/sql 发起标准 SQL 请求 ▼ Pinot Broker(标准 SQL 查询端点,即 PQL 端点弃用后的替代方案)

几个值得注意的实现事实:

  • PinotDbApiHook明确注释:使用标准 SQL 端点,因为 PQL 端点即将被弃用(对应官方查询文档);conn_type = "pinot"hook_name = "Pinot Broker",这些注册信息同时出现在 get_provider_info.py 与 provider.yaml 中;
  • 连接信息中的login/password会作为username/password传给pinotdb.connect,可用于带认证的 Broker;
  • SQLExecuteQueryOperator.execute()do_xcom_pushrequires_result_fetch为真时才传入handler抓取结果,否则查询结果不会显式拉取,这一点对超大结果集的查询可以起到节省内存的作用。

延伸:用 Hook 完成离线数据管理

虽然本文主题是查询算子,但 Pinot Provider 还提供PinotAdminHook,用于调用pinot-admin.sh脚本完成离线数据摄入(AddSchema、AddTable、CreateSegment、UploadSegment 四个子命令),可作为构建"离线灌数 + 在线查询"完整链路时的补充手段,详见 hooks.rst。PinotAdminHook的关键行为包括:

  • 在 4.0.0 版本起cmd_path被硬编码为pinot-admin.sh,必须确保该脚本在 PATH 中,传入其他值会直接抛出RuntimeError
  • 由于早期 Pinot 的pinot-admin.sh无论成败都返回退出码 0,可通过pinot_admin_system_exit标志(或连接Extra中的同名键)切换为"按输出内容判断":当输出中包含ErrorException时视为失败并抛出AirflowException

结语

Apache Airflow 与 Apache Pinot 的集成遵循"通用 SQL 优先"的设计哲学:没有专属 Operator,不代表能力缺失。借助SQLExecuteQueryOperator+PinotDbApiHook,你可以在一个统一、可模板化、支持 XCom 传递结果的框架下,对 Pinot 执行标准 SQL 聚合查询。实际操作时请记住三件事:安装apache-airflow-providers-apache-pinot并确保pinotdb依赖可用、按 Broker 实际情况配置pinot类型 Connection(Host/Port/Extra 端点)、利用算子显式参数覆盖连接默认值。系统测试示例 example_pinot.py 与 example_pinot_dag.py 是开箱即用的参考实现,可直接作为新 DAG 的起点。

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

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 15:47:54

手写文字擦除方案拆解:从模型结构到数据管线的工程实践

简介&#xff1a;面向大学生竞赛与深度学习实践者的手写文字擦除一等奖方案完整资源包。该方案针对试卷扫描图中手写红黑蓝笔迹、手画线段、污渍脏点与印刷字重叠等复杂场景&#xff0c;提供从数据划分、模型训练到测试推理的全流程实现。官方训练集共1081对&#xff0c;方案另…

作者头像 李华
网站建设 2026/9/13 15:44:47

Python资产管理系统部署实战:从解压到扩展功能

简介&#xff1a;这是一份基于Python开发的资产管理系统源码包&#xff0c;面向企业或个人对硬件设备、软件资源进行登记、跟踪与维护的场景&#xff0c;适合正在学习Python Web开发、想通过完整项目提升实战能力的中初级开发者。该压缩包共包含41个文件&#xff0c;以21个Pyth…

作者头像 李华
网站建设 2026/9/13 15:44:04

多目标推荐系统实战:Otto竞赛LightGBM单模型0.594技术拆解

简介&#xff1a;对标Kaggle Otto多目标推荐系统赛题&#xff0c;这是一份单模型LB分数0.594、排名约30的完整源代码方案&#xff0c;适合想冲击推荐类竞赛榜单的选手及希望深入多目标推荐工程的数据科学学习者。代码覆盖数据处理、用户/物品特征与相似度特征构建、协同过滤与图…

作者头像 李华
网站建设 2026/9/13 15:43:36

Python+OpenCV实现相机标定:从棋盘格到内参矩阵的完整方案

简介&#xff1a;面向计算机视觉初学者与机器人、三维重建等领域的开发者&#xff0c;这份相机标定Python程序提供了完整可用的内参求解方案。资源自带1110规格的正友棋盘格图片&#xff0c;既可打印后拍摄&#xff0c;也可直接放在显示器上配合附带程序使用&#xff0c;大幅降…

作者头像 李华