news 2026/10/9 12:10:50

在 Airflow DAG 中编排 Hive:apache-airflow-providers-apache-hive 9.3.0 实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
在 Airflow DAG 中编排 Hive:apache-airflow-providers-apache-hive 9.3.0 实战指南

【免费下载链接】context-hub

项目地址:https://gitcode.com/gh_mirrors/co/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 / BeelineHiveOperator、HiveCliHook在 Airflow 任务中通过 Hive CLI 或 Beeline 执行 HQL
Hive MetastoreHivePartitionSensor及 metastore hooks等待分区出现、通过 metastore 检查表元数据
HiveServer2HiveServer2Hook从 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

项目地址:https://gitcode.com/gh_mirrors/co/context-hub
点击查看免费下载
上一篇:komorebi 的 container-padding 命令:按工作区精确控制容器内边距
下一篇:brpc ExecutionQueue 深入指南:wait-free 异步串行任务队列的设计与实战

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

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

机器学习预测钢管混凝土柱承载力:高精度代理模型实战

简介&#xff1a;本资源是一套面向土木工程与人工智能交叉领域研究者的机器学习建模实践项目&#xff0c;聚焦于内配型钢钢管混凝土柱承载力的高精度预测问题&#xff0c;适用于结构工程方向的研究生、科研人员及具备Python基础的算法实践者。压缩包共5个文件&#xff0c;含4个…

作者头像 李华
网站建设 2026/10/9 12:09:12

VS Code 国际化插件 i18n Ally 配置到 TaoToken 的完整实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 12:07:36

ARIMA预测新能源汽车销量:从数据准备到滚动验证的完整指南

简介&#xff1a;基于ARIMA模型的新能源汽车销量预测PDF&#xff0c;是一份面向汽车行业数据分析人员、高校研究者和市场预测从业者的时间序列建模参考。资源为单个PDF文档&#xff0c;大小约1.11MB&#xff0c;收录了完整的期刊论文内容&#xff0c;详细展示ARIMA模型应用流程…

作者头像 李华