尧图精选

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

🕒 发布时间:2026/9/13 15:49:40 📁 来源:尧图网络
Apache Airflow 集成 Apache Pinot 实战使用 SQLExecuteQueryOperator 执行实时 OLAP 查询【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 官方并没有为 Apache Pinot 提供专门的 Operator而是统一推荐使用通用 SQL 执行算子SQLExecuteQueryOperator来对 Pinot Broker 发起标准 SQL 查询。本文以仓库中的 operators.rst 为核心完整讲解连接配置、示例 DAG、底层 Hook 调用链与参数优先级帮助你快速在 Airflow 中构建针对 Pinot 的查询任务并掌握可复用的最佳实践。为什么 Pinot 没有专属 OperatorApache Pinot 是一个面向列的、分布式的开源 OLAP 数据存储由 Java 编写专为低延迟分析场景设计适合在不可变数据上进行聚合类快速分析也支持实时数据摄入。Pinot 对外暴露的是标准 SQL 查询接口因此 Airflow 社区将其接入方式统一收敛到通用 SQL 体系使用SQLExecuteQueryOperator位于airflow.providers.common.sql.operators.sql执行查询底层由 Pinot Provider 提供的PinotDbApiHook负责建立与 Pinot Broker 的连接。这一点在文档中有两条明确的note说明Apache Pinot 没有专属 Operator请直接使用SQLExecuteQueryOperator必须先安装对应的 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-airflow2.11.0apache-airflow-providers-common-compat1.10.1apache-airflow-providers-common-sql1.32.0pinotdb5.1.0其中pinotdb是 Pinot 官方的 Python DB-API 驱动PinotDbApiHook正是通过它连接 Broker。配置 Pinot 连接Connection使用SQLExecuteQueryOperator时通过conn_id参数指向一个名为pinot类型的 Airflow Connection。连接元数据的结构如下参数取值Host: stringPinot Broker 的主机名或 IP 地址Port: intPinot Broker 端口默认8000Schema: string不使用Extra: JSON可选字段例如{endpoint: query/sql}注意事项文档表格中默认端口写为 8000但实际生产环境常见的 Broker 端口是 8099 或 9000 等请以你实际部署的 Pinot Broker 端口为准Schema字段在查询场景下不被使用但PinotDbApiHook.get_conn()会读取extra_dejson中的schema键作为请求协议scheme默认httpExtra中的endpoint用于指定 Broker 的 SQL 查询端点默认值为/query/sql。更深入的实现细节可以参考 pinot.py 中PinotDbApiHook.get_conn()的源码它通过pinotdb.connect(host..., port..., usernameconn.login, passwordconn.password, path..., scheme...)建立连接。也就是说Airflow Connection 的Login与Password字段会被透传给 Pinot 做身份认证Extra中支持的键包括endpointSQL 查询端点路径默认query/sqlschema连接协议默认http对应conn_type。此外PinotDbApiHook.get_uri()会拼出形如http://localhost:9000/query/sql的 URIconn_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_idDAG_ID, start_datedatetime.datetime(2025, 1, 1), default_args{conn_id: my_pinot_conn}, scheduleonce, catchupFalse, ) as dag: # Task: Simple query to test connection and query engine select_1_task SQLExecuteQueryOperator( task_idselect_1, sqlSELECT 1, ) # Task: Count total records in airlineStats (sample table) count_airline_stats SQLExecuteQueryOperator( task_idcount_airline_stats, sqlSELECT COUNT(*) FROM airlineStats, ) # Task: Group by Carrier and count flights group_by_carrier SQLExecuteQueryOperator( task_idgroup_by_carrier, sqldedent( 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代码要点拆解conn_id 的两种指定方式示例在 DAG 的default_args中统一设置了conn_id: my_pinot_conn因此每个SQLExecuteQueryOperator无需重复传参你也可以在单个 Task 上通过conn_id显式覆盖。sql参数支持多行 SQLSQLExecuteQueryOperator的template_fields (sql, parameters, ...)说明sql是支持模板渲染的字段。多行语句可以用textwrap.dedent或三引号书写保持可读性。任务编排示例用位运算将三个查询串成线性链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继承自 BaseOperator为True时结果会自动写入 XCom供下游任务消费注意PinotDbApiHook声明了supports_autocommit False且其set_autocommit、insert_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_idpinot 类型 → 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可用于带认证的 BrokerSQLExecuteQueryOperator.execute()在do_xcom_push或requires_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中的同名键切换为按输出内容判断当输出中包含Error或Exception时视为失败并抛出AirflowException。结语Apache Airflow 与 Apache Pinot 的集成遵循通用 SQL 优先的设计哲学没有专属 Operator不代表能力缺失。借助SQLExecuteQueryOperatorPinotDbApiHook你可以在一个统一、可模板化、支持 XCom 传递结果的框架下对 Pinot 执行标准 SQL 聚合查询。实际操作时请记住三件事安装apache-airflow-providers-apache-pinot并确保pinotdb依赖可用、按 Broker 实际情况配置pinot类型 ConnectionHost/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),仅供参考
上一篇/下一篇内容由系统自动关联 返回资讯列表 →