【免费下载链接】context-hub
本指南围绕 Apache Airflow 官方 Hive provider(apache-airflow-providers-apache-hive9.3.0)展开,讲解如何在 DAG 中运行 HQL、等待分区就绪、读取 Hive 元数据以及通过 HiveServer2 从 Python 任务代码查询数据。读完本文,你将掌握该 provider 的安装方式、三类连接(hive_cli/hiveserver2/hive_metastore)的配置方法,以及HiveOperator、HivePartitionSensor、HiveMetastoreHook、HiveServer2Hook的完整用法,能够直接在现有 Airflow 项目中落地 Hive 相关任务。
本文依据仓库中的文档 providers-apache-hive/python/DOC.md 编写,该文档对应 provider 版本 9.3.0,内容由维护者提供并归档在 Context Hub 的 Apache Airflow 文档集中。
何时使用这个 Provider
apache-airflow-providers-apache-hive用于Airflow DAG 需要执行 Hive SQL、等待 Hive 分区出现、或从 Python 任务代码读取 Hive 元数据的场景。
需要特别强调的一点是:这是一个 Airflow provider,不是面向普通 Python 应用的独立 Hive 客户端。如果你在 Airflow 之外编写常规 Python 程序,应直接使用专门的 Hive 客户端库,而不是引入 Airflow 的 hooks 与 operators。只有当 Airflow 需要把 Hive 工作编排为 DAG 任务时,才应当选用该 provider。
安装与版本配套
将 provider 安装到与apache-airflow相同的 Python 环境或容器镜像中。实际操作上,这意味着scheduler、webserver 以及每一个导入 DAG 代码的 worker 都必须具备该 provider——只要有一个 worker 镜像缺失,任务执行阶段就会出现ModuleNotFoundError。
安装时应将 Airflow 与 provider 版本一起固定,并使用与你的 Airflow 版本对应的 constraints 文件:
AIRFLOW_VERSION="<your-airflow-version>" PROVIDER_VERSION="9.3.0" PYTHON_VERSION="$(python -c 'import sys; print(f"{sys.version_info.major}.{sys.version_info.minor}")')" CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt" python -m pip install \ "apache-airflow==${AIRFLOW_VERSION}" \ "apache-airflow-providers-apache-hive==${PROVIDER_VERSION}" \ --constraint "${CONSTRAINT_URL}"使用 constraints 文件可以保证 provider 及其传递依赖与当前 Airflow 版本兼容,避免因依赖版本漂移导致的意外行为。请将AIRFLOW_VERSION替换为你的实际 Airflow 版本号。
选择正确的接口:三条集成路径
该 provider 对外暴露三条常见的集成路径,分别对应不同的使用场景:
| 集成路径 | 核心组件 | 适用场景 |
|---|---|---|
| Hive CLI / Beeline | HiveOperator、HiveCliHook | 在 Airflow 任务中通过 Hive CLI 或 Beeline 执行 HQL |
| Hive Metastore | HivePartitionSensor及 metastore hooks | 等待分区出现、通过 metastore 检查表元数据 |
| HiveServer2 | HiveServer2Hook | 从 Python 任务代码通过 HiveServer2 查询 Hive |
大多数 DAG 只用到下面三种典型模式:
- 用
HiveOperator运行 HQL; - 用
HivePartitionSensor等待分区; - 用
HiveMetastoreHook检查表状态。
配置 Airflow Connections
推荐的实践是:把主机名、用户名、Kerberos 细节和 SSL 设置放在 Airflow connection 中,而不是硬编码在 DAG 文件里。
该 provider 使用的典型 connection id 有:
hive_cli_default— 供HiveOperator和HiveCliHook使用;hiveserver2_default— 供HiveServer2Hook使用;metastore_default— 供 metastore hooks 和分区传感器使用。
下面是通过环境变量配置这三种 connection 的示例:
export AIRFLOW_CONN_HIVE_CLI_DEFAULT='{"conn_type":"hive_cli","host":"hs2.example.com","port":10000,"login":"airflow","schema":"default","extra":{"use_beeline":true}}' export AIRFLOW_CONN_HIVESERVER2_DEFAULT='{"conn_type":"hiveserver2","host":"hs2.example.com","port":10000,"login":"airflow","schema":"default"}' export AIRFLOW_CONN_METASTORE_DEFAULT='{"conn_type":"hive_metastore","host":"metastore.example.com","port":9083}'关键字段说明:
conn_type— 必须与各组件期望的类型一致:hive_cli、hiveserver2、hive_metastore;host/port— HiveServer2 默认端口为10000,Hive Metastore 默认端口为9083;login— 连接使用的用户名;schema— 默认数据库(schema),例如default或analytics;extra— 存放额外选项,如{"use_beeline": true}表示使用 Beeline 而非 Hive CLI。
如果你的集群使用了Kerberos、SSL、LDAP 或自定义 Beeline 选项,请把这些值放进 connection 的 extras 中,而不是嵌入 Python 代码。认证信息与集群专属配置同样应放在 Airflow connections 或你的 secrets backend 中,而不是写在 DAG 源文件里。
在 DAG 中运行 Hive SQL
当任务需要在 DAG 执行过程中运行 HQL 时,使用HiveOperator。它通过配置的hive_cli_conn_id所指向的 connection(默认hive_cli_default)调用 worker 本地的 Hive CLI 或 Beeline。
from airflow import DAG from airflow.providers.apache.hive.operators.hive import HiveOperator from pendulum import datetime with DAG( dag_id="hive_load_daily_partition", start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False, ) as dag: create_table = HiveOperator( task_id="create_table", hive_cli_conn_id="hive_cli_default", hql=""" CREATE TABLE IF NOT EXISTS analytics.events ( user_id STRING, event_name STRING ) PARTITIONED BY (ds STRING) STORED AS PARQUET """, ) load_partition = HiveOperator( task_id="load_partition", hive_cli_conn_id="hive_cli_default", hql=""" INSERT OVERWRITE TABLE analytics.events PARTITION (ds='${hiveconf:ds}') SELECT user_id, event_name FROM staging.events_raw WHERE ds='${hiveconf:ds}' """, hiveconfs={"ds": "{{ ds }}"}, ) create_table >> load_partition要点解析:
hql参数传入要执行的 Hive 查询文本,支持多行字符串;- 当你想让 Airflow 的模板引擎把值注入 HQL,而不在 Python 里用字符串拼接 SQL时,使用
hiveconfs参数:hiveconfs={"ds": "{{ ds }}"}会把 Airflow 的执行日期{{ ds }}作为 Hive 配置变量hiveconf:ds传入,HQL 中通过${hiveconf:ds}引用; - 这个例子体现了典型的"建表 → 按天装载分区"工作流,两条任务通过
create_table >> load_partition建立依赖关系。
注意:HiveOperator和HiveCliHook在 worker 上运行的是本地 Hive CLI 或 Beeline,因此二进制程序必须存在于 worker 镜像中并位于PATH上。如果使用 Beeline,请在 connection 中设置启用 Beeline,同时确保 worker 具备 JDBC 驱动以及所需的集群客户端配置。
等待一个分区出现
当下游任务需要等到 metastore 报告某个特定分区存在后才能继续时,使用HivePartitionSensor。典型场景是:另一个系统负责发布 Hive 分区,你的 DAG 只有在分区在 metastore 中可见之后才能继续执行。
from airflow import DAG from airflow.providers.apache.hive.sensors.hive_partition import HivePartitionSensor from pendulum import datetime with DAG( dag_id="wait_for_hive_partition", start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False, ) as dag: wait_for_partition = HivePartitionSensor( task_id="wait_for_partition", table="analytics.events", partition="ds='{{ ds }}'", metastore_conn_id="metastore_default", poke_interval=60, timeout=60 * 60, )参数说明:
table— 要检查的表名,如analytics.events;partition— 期望的分区表达式,如"ds='{{ ds }}'",同样支持 Airflow 模板;metastore_conn_id— metastore 连接 id(默认metastore_default);poke_interval— 轮询间隔(秒),示例为 60 秒;timeout— 超时时间(秒),示例为60 * 60,即 1 小时;超时后传感器会失败并让任务重试或失败。
从 Python 任务读取元数据
当任务代码需要检查表或分区(而不只是等待它们出现)时,使用HiveMetastoreHook。这类用法非常适合分支逻辑、审计或启动较重下游工作前的健全性检查。
from airflow import DAG from airflow.decorators import task from airflow.providers.apache.hive.hooks.hive import HiveMetastoreHook from pendulum import datetime with DAG( dag_id="inspect_hive_partitions", start_date=datetime(2026, 1, 1), schedule=None, catchup=False, ) as dag: @task def print_latest_partition() -> str | None: hook = HiveMetastoreHook(metastore_conn_id="metastore_default") latest = hook.max_partition("analytics", "events", field="ds") print(f"latest partition: {latest}") return latest print_latest_partition()示例中通过HiveMetastoreHook.max_partition(table, partition, field=...)查询指定分区字段的最大值。返回的latest既可以直接打印,也可以作为 Python 任务函数的返回值传给后续任务,支撑分支决策。
注意:HivePartitionSensor和HiveMetastoreHook连接的是 Hive metastore,而不是 HiveServer2。一个能正常工作的查询连接并不能保证 metastore 连接也是正确的,两类连接需要分别验证。
通过 HiveServer2 从 Python 查询数据
当 Python 任务需要拉取 Hive 中的行数据(而不是通过 CLI 提交 HQL)时,使用HiveServer2Hook。这在需要以查询结果驱动控制流决策的场景下尤其有用。
from airflow import DAG from airflow.decorators import task from airflow.providers.apache.hive.hooks.hive import HiveServer2Hook from pendulum import datetime with DAG( dag_id="query_hiveserver2", start_date=datetime(2026, 1, 1), schedule=None, catchup=False, ) as dag: @task def read_counts() -> None: hook = HiveServer2Hook( hiveserver2_conn_id="hiveserver2_default", schema="analytics", ) rows = hook.get_records( """ SELECT ds, COUNT(*) AS row_count FROM events GROUP BY ds ORDER BY ds DESC LIMIT 7 """ ) for ds, row_count in rows: print(ds, row_count) read_counts()要点解析:
HiveServer2Hook在构造时接收hiveserver2_conn_id(默认hiveserver2_default)和可选的schema参数,用于切换默认数据库;get_records(sql)执行查询并返回行记录,可以直接在 Python 中迭代处理;- 该方式适合小结果集与控制流决策。对于大规模数据移动,应把工作保留在 Hive SQL 任务中,而不是通过 Python worker 拉取大结果集——这会带来不必要的网络与内存开销。
关键注意事项汇总
- Worker 镜像必须包含 CLI 二进制:
HiveOperator和HiveCliHook运行 worker 上的本地 Hive CLI 或 Beeline,二进制需存在于 worker 镜像并位于PATH;若使用 Beeline,还需要 JDBC 驱动与集群客户端配置; - Metastore 与 HiveServer2 是两条不同的通道:
HivePartitionSensor、HiveMetastoreHook走 metastore(默认端口 9083),HiveServer2Hook走 HiveServer2(默认端口 10000),查询连通不代表 metastore 连通; - 在每一个导入 DAG 代码的位置安装 provider:一个缺失的 worker 镜像就足以在任务执行时引发
ModuleNotFoundError; - 凭据与集群专属配置放入 connections / secrets backend:不要在 DAG 源文件中硬编码认证信息;
- 使用场景有边界:provider 面向 Airflow 编排 Hive 工作;Airflow 之外的普通 Python 应用应使用专门的 Hive 客户端。
总结
apache-airflow-providers-apache-hive(9.3.0)为 Airflow DAG 提供了三条与 Hive 交互的通道:通过HiveOperator/HiveCliHook执行 HQL、通过HivePartitionSensor/HiveMetastoreHook与 metastore 交互、通过HiveServer2Hook从 Python 任务查询数据。配合hive_cli_default、hiveserver2_default、metastore_default三种 connection 的组织方式,即可在不把连接细节硬编码进 DAG 的前提下,把 Hive 的建表、装载、分区等待与元数据检查完整地编排进 Airflow 工作流。更完整的组件 API 说明可参阅仓库中的 DOC.md 原文,以及 Context Hub 内容指南 了解文档组织方式。
【免费下载链接】context-hub
相关推荐
在 Apache Airflow DAG 中编排 Apache Beam 管道:apache-airflow-providers-apache-beam 实战指南
在 Apache Airflow DAG 中编排 Apache Beam 管道:apache airflow providers apache beam 实战指
在 Airflow DAG 中运行 Pig Latin:apache-airflow-providers-apache-pig 4.8.2 实战指南
在 Airflow DAG 中运行 Pig Latin:apache airflow providers apache pig 4.8.2 实战指南 apach
Airflow DAG 中提交 Spark 应用:apache-airflow-providers-apache-spark 5.5.1 实战指南
Airflow DAG 中提交 Spark 应用:apache airflow providers apache spark 5.5.1 实战指南 本篇技术指南
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考