
基于 Apache Airflow 的 Apache Drill 提供商实战指南安装配置、连接管理、SQL 查询执行与 DrillHook 源码解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow# Apache Airflow Apache Drill 提供商实战指南连接配置、SQL 查询执行与 DrillHook 源码解析本篇技术指南以 Apache Airflow 官方仓库中apache-airflow-providers-apache-drill提供商当前版本 3.3.3为主线系统讲解如何在 Airflow 中接入 Apache Drill 查询引擎从包的安装与版本依赖、Drill 连接Connection的配置要点到使用SQLExecuteQueryOperator编排 Drill SQL 任务再到深入DrillHook源码理解底层连接串的构造逻辑与 Drill 引擎的语义限制。读完本文你将能够独立完成 Drill 提供商的环境搭建、连接配置并编写出可复用、可上生产的 Drill 查询 DAG。一、提供商包概览apache-airflow-providers-apache-drill是 Apache Airflow 面向 Apache Drill 中的元数据声明该提供商当前状态为ready、生命周期为production属于可直接用于生产环境的稳定模块。它提供的核心能力包括一个数据库连接类型drillconnection-type: drillhook 名Drill用于在 Airflow 中登记 Drillbit 连接信息DrillHookairflow.providers.apache.drill.hooks.drill.DrillHook基于 sqlalchemy-drill 与 Drill 引擎交互通过通用 SQL 运算符SQLExecuteQueryOperator执行 Drill SQL 查询的能力。提供商包按 Airflow 惯例提供完整的文档树其核心文档入口为 providers/apache/drill/docs/index.rst与连接、运算符相关的详细指南分别位于 连接指南 与 运算符指南。二、安装与版本要求2.1 安装命令在已有 Airflow 环境之上安装该提供商包只需一条 pip 命令pip install apache-airflow-providers-apache-drill根据 README.rst该包支持 Python 3.10、3.11、3.12、3.13、3.14 版本。2.2 依赖清单包安装时会对依赖做版本校验最低依赖要求如下出自 index.rst 与 README.rstPIP 包版本要求apache-airflow2.11.0apache-airflow-providers-common-sql1.32.0apache-airflow-providers-common-compat1.8.0sqlalchemy-drill1.1.0,!1.1.6,!1.1.7其中sqlalchemy-drill是关键的第三方桥接库它实现了 SQLAlchemy 方言与驱动使 Airflow 能够通过 SQLAlchemy 引擎访问 Drillbit。版本约束排除了1.1.6与1.1.7两个版本安装时需注意不要锁定到被排除的版本。该包还依赖common-sql与common-compat两个基础提供商前者提供了DbApiHook与SQLExecuteQueryOperator等通用 SQL 抽象这正是本提供商只提供 Hook、复用通用运算符这一设计的基础。2.3 官方发布包校验正式发布的包可从 Apache 官方下载站获取 sdist 与 wheel 两种格式并配套提供.ascPGP 签名与.sha512哈希校验文件例如apache_airflow_providers_apache_drill-3.3.3.tar.gz及对应的校验文件。下载后建议先核对 sha512 与签名确保包完整性。三、配置 Apache Drill 连接Connection在 Airflow 中执行 Drill 查询之前必须先配置一个 Drill 类型的连接。详细配置说明见 连接指南。3.1 默认连接 IDDrillHook与相关运算符默认使用连接 IDdrill_default。在源码中可以看到这一约定conn_name_attr drill_conn_id default_conn_name drill_default conn_type drill hook_name Drill见 drill.py也就是说如果在 Airflow UI 的 Admin → Connections 中新建一个连接把 Conn Id 填为drill_default、Conn Type 选为Drill那么使用该连接的所有运算符都无需再显式传入conn_id。3.2 连接字段说明配置 Drill 连接时需要关注以下字段Host必填连接目标 Drillbit 的主机地址。根据连接方式不同含义略有区别HTTP / JDBC 方式填写 Drillbit 的主机名或 IPODBC 方式填写 Drill ODBC 连接的 DSN。Port可选Drillbit 的端口。Drill 默认 REST API 端口为8047若不填写DrillHook.get_conn()会拼接出形如host:的空端口形式因此实践中建议显式配置。Extra可选一个 JSON 字典用于传递 sqlalchemy-drill 支持的额外参数。该提供商支持两个关键键Extra 键含义默认值dialect_driversqlalchemy-drill 使用的方言与驱动标识drill_sadrill即 HTTP REST 方式storage_plugin该连接默认使用的 Drill 存储插件dfs典型配置示例{ dialect_driver: drill_sadrill, storage_plugin: dfs }在 UI 的 Extra 文本框中直接粘贴上述 JSON 即可。3.3 底层连接串构造从源码看配置如何生效连接配置最终如何变成一条可执行的数据库 URL关键在于DrillHook.get_conn()的实现见 drill.pydef get_conn(self) - PoolProxiedConnection: Establish a connection to Drillbit. conn_md self.get_connection(self.get_conn_id()) creds f{conn_md.login}:{conn_md.password} if conn_md.login else database_url ( f{conn_md.extra_dejson.get(dialect_driver, drillsadrill)}://{creds} f{conn_md.host}:{conn_md.port}/ f{conn_md.extra_dejson.get(storage_plugin, dfs)} ) if ? in database_url: raise ValueError(Drill database_url should not contain a ?) engine create_engine(database_url) ... return engine.raw_connection()可以拆解出以下几点关键逻辑URL 拼接规则最终 URL 形如drillsadrill://user:passhost:port/dfs。其中dialect_driver作为 SQLAlchemy 方言前缀源码默认取drillsadrill即 HTTP REST 方言若连接中配置了login则自动拼入login:password凭据段storage_plugin作为 URL 的路径部分默认dfs。?校验如果拼接后的 URL 中包含?直接抛出ValueError(Drill database_url should not contain a ?)——因为 Drill 的数据库 URL 不允许携带查询字符串形式的参数。引擎创建使用sqlalchemy.create_engine创建引擎并返回raw_connection()后续 SQL 执行全部经由 SQLAlchemy 会话完成。注意连接指南中 Extra 的dialect_driver默认值写作drill_sadrill而源码get_conn()中实际采用的默认字符串是drillsadrill——两者是同一 HTTP 方言在文档表述与SQLAlchemy 方言标识上的两种写法配置时两者皆可被 sqlalchemy-drill 识别。3.4 get_uri()面向 Airflow 生态的连接 URI 输出DrillHook还实现了get_uri()见 drill.py将连接信息序列化为标准 URI例如drill://localhost:8047/dfs?dialect_driverdrillsadrill。该 URI 常用于 Airflow 中需要展示或传递连接串的场景如 OpenLineage 等元数据系统。单元测试 test_drill.py 通过 7 组参数化用例验证了get_uri()的边界行为从中可以看到不配置端口时输出drill://host/dfs?dialect_driverdrillsadrill省略端口段不配置任何 Extra 字段时全部回落到默认值dfs与drillsadrill自定义conn_type、dialect_driver、storage_plugin时URI 相应变化如custom://myhost:1234/myplugin?dialect_drivermydriver。四、在 DAG 中执行 Drill SQLSQLExecuteQueryOperator4.1 推荐用法SQLExecuteQueryOperator在 DAG 中向 Drill 提交 SQL官方推荐使用通用 SQL 运算符SQLExecuteQueryOperatorairflow.providers.common.sql.operators.SQLExecuteQueryOperator。它能够在一个 Drillbit 上执行一条或多条 SQL 语句支持sql参数模板化Jinja 渲染支持将 SQL 放到外部.sql文件中引用。使用前需要按第三节配置好 Drill 连接drill_default或自定义 ID通过conn_id参数把连接传给运算符。官方系统测试示例 DAG example_drill_dag.py 给出了完整可运行的最小示例from datetime import datetime from airflow.models import DAG from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator with DAG( dag_idexample_drill_dag, scheduleNone, start_datedatetime(2021, 1, 1), catchupFalse, tags[example], ) as dag: sql_task SQLExecuteQueryOperator( task_idjson_to_parquet_table, sql drop table if exists dfs.tmp.employee; create table dfs.tmp.employee as select * from cp.employee.json; , )这个例子展示了 Drill 的一个典型用法把 JSON 数据文件转换为 Parquet 表——cp.是 Drill 自带的 classpath 存储插件内置示例数据employee.jsondfs.tmp.是文件系统存储插件下的临时工作区。两条 SQL 先删除旧表再通过create table ... as selectCTAS完成格式转换非常适合作为 Drill 数据管道中的 ETL 步骤。提示sql参数支持模板化你可以用{{ ds }}等 Airflow 模板变量动态拼接查询或将 SQL 提取到.sql文件后通过相对路径引用便于维护与复用。4.2 弃用警告请勿再使用 DrillOperator需要特别提醒历史版本中执行 Drill SQL 使用的是DrillOperator该运算符目前已弃用deprecated并将在未来版本中移除。运算符指南 operators.rst 中对此给出了明确警告当前源码中operators/__init__.py已不再导出任何运算符类见 operators/init.py运算符能力完全由 common-sql 的SQLExecuteQueryOperator承担。迁移建议如果现有 DAG 还在使用DrillOperator应尽快将task_id、sql、conn_id等参数平移到SQLExecuteQueryOperator参数语义基本一致以规避未来升级时的破坏性变更。五、DrillHook 源码解析与 Drill 引擎交互的底层实现5.1 类结构DrillHook继承自airflow.providers.common.sql.hooks.sql.DbApiHook见 drill.py因此天然继承了 common-sql 提供的一整套数据库 APIrun(sql)执行 SQLget_first(sql)取首行结果get_records(sql)取全部记录get_df(sql, df_type...)结果转为 DataFrame支持pandas与polars两种类型get_pandas_df、get_df_by_chunks等衍生方法。继承带来的另一大收益是钩子级数据血缘Hook Level Lineage自提供商 3.3.0 版本起见 changelog.rst 中的 feat: Add Hook Level Lineage to SQL hooksrun与get_df等方法会自动调用send_sql_hook_lineage上报 SQL 血缘信息。单元测试 test_drill.py 中test_run_hook_lineage、test_get_df_hook_lineage、test_get_df_by_chunks_hook_lineage三个用例分别验证了run、get_df、get_df_by_chunks路径下的血缘上报行为。5.2 Drill 引擎语义限制无事务、无 INSERTDrillHook针对 Drill 引擎的语义特性覆写了两个在传统关系型数据库上常用、但在 Drill 中不成立的方法均直接抛出NotImplementedErrordef set_autocommit(self, conn: Connection, autocommit: bool) - NoReturn: raise NotImplementedError(There are no transactions in Drill.) def insert_rows(self, table, rows, target_fieldsNone, commit_every1000, replaceFalse, **kwargs): raise NotImplementedError(There is no INSERT statement in Drill.)见 drill.py这两处覆写背后是 Drill 引擎的两个核心特性Drill 不提供事务因此supports_autocommit False任何set_autocommit调用都会报错Drill 没有 INSERT 语句写入数据只能通过 CTASCREATE TABLE AS SELECT、COPY INTO或INSERT INTO ... SELECT部分版本等批量方式完成insert_rows自然没有实现。单元测试 test_drill.py 中的test_set_autocommit_raises_not_implemented与test_insert_rows_raises_not_implemented分别验证了这两个行为。对 DAG 编写者的启示在 Drill 任务中不要依赖事务语义也不要用常规的逐行 INSERT 模式数据写入应使用 CTAS / COPY INTO 等批量 SQL如第四节示例中create table ... as select的写法。5.3 连接串安全校验get_conn()中的?校验同样有对应的单元测试test_drill.py当 Host 中包含?时host_with?get_conn()会抛出ValueError(Drill database_url should not contain a ?)正常主机名则成功建立连接。这提醒我们在配置 Drill 连接时Host 字段务必只填主机名/IP不要混入协议前缀或查询参数。六、测试体系与验证方式该提供商围绕 Hook 提供了完整的测试覆盖单元测试providers/apache/drill/tests/unit/apache/drill/hooks/test_drill.py —— 覆盖连接建立、URI 序列化、get_first/get_records/get_dfpandas 与 polars、无事务/无 INSERT 异常、血缘上报等系统测试providers/apache/drill/tests/system/apache/drill/example_drill_dag.py —— 即上文第四节展示的示例 DAG需要在真实 Drill 环境中运行集成测试providers/apache/drill/tests/integration/apache/drill/hooks/test_drill.py —— 面向真实 Drillbit 的连接与查询验证。对于想要在本地验证的开发者可以先启动一个 Drillbit默认 REST 端口 8047配置好drill_default连接然后运行系统测试 DAG 或直接执行DrillHook().run(select * from cp.employee.jsonlimit 5)做冒烟验证。七、总结apache-airflow-providers-apache-drill是 Airflow 接入 Apache Drill 的官方生产级通道其设计遵循了 Airflow 3.x 时代提供商只提供 Hook 与连接类型、运算符复用 common-sql的架构方向连接层通过drill连接类型 DrillHook把 Host/Port/Extradialect_driver、storage_plugin映射为drillsadrill://.../dfs形式的 SQLAlchemy URL执行层用SQLExecuteQueryOperator提交可模板化的 SQL实现 JSON→Parquet 等典型 ETL 场景语义适配层明确声明 Drill 无事务、无 INSERT指导用户使用 CTAS 等批量写入生态集成层天然获得 lineage 血缘上报、DataFramepandas/polars结果处理等能力。配合本仓库中的 连接指南、运算符指南、changelog.rst 与 源码你可以在此基础上继续深入探索 Drill 的存储插件、CTAS 分区写入等高级用法构建稳定高效的 Drill 数据工作流。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考