首页 / 资讯中心 / 文章详情

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

在 Airflow DAG 中编排 Hive:apache-airflow-providers-apache-hive 9.3.0 实战指南 ★ FEATURED ARTICLE
【免费下载链接】context-hub项目地址https://gitcode.com/gh_mirrors/co/context-hub点击查看免费下载本指南围绕 Apache Airflow 官方 Hive providerapache-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 文档集中。何时使用这个 Providerapache-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_VERSIONyour-airflow-version PROVIDER_VERSION9.3.0 PYTHON_VERSION$(python -c import sys; print(f{sys.version_info.major}.{sys.version_info.minor})) CONSTRAINT_URLhttps://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 执行 HQLHive 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_metastorehost/port— HiveServer2 默认端口为10000Hive Metastore 默认端口为9083login— 连接使用的用户名schema— 默认数据库schema例如default或analyticsextra— 存放额外选项如{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_idhive_load_daily_partition, start_datedatetime(2026, 1, 1), scheduledaily, catchupFalse, ) as dag: create_table HiveOperator( task_idcreate_table, hive_cli_conn_idhive_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_idload_partition, hive_cli_conn_idhive_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_idwait_for_hive_partition, start_datedatetime(2026, 1, 1), scheduledaily, catchupFalse, ) as dag: wait_for_partition HivePartitionSensor( task_idwait_for_partition, tableanalytics.events, partitionds{{ ds }}, metastore_conn_idmetastore_default, poke_interval60, timeout60 * 60, )参数说明table— 要检查的表名如analytics.eventspartition— 期望的分区表达式如ds{{ ds }}同样支持 Airflow 模板metastore_conn_id— metastore 连接 id默认metastore_defaultpoke_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_idinspect_hive_partitions, start_datedatetime(2026, 1, 1), scheduleNone, catchupFalse, ) as dag: task def print_latest_partition() - str | None: hook HiveMetastoreHook(metastore_conn_idmetastore_default) latest hook.max_partition(analytics, events, fieldds) print(flatest 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_idquery_hiveserver2, start_datedatetime(2026, 1, 1), scheduleNone, catchupFalse, ) as dag: task def read_counts() - None: hook HiveServer2Hook( hiveserver2_conn_idhiveserver2_default, schemaanalytics, ) 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默认端口 9083HiveServer2Hook走 HiveServer2默认端口 10000查询连通不代表 metastore 连通在每一个导入 DAG 代码的位置安装 provider一个缺失的 worker 镜像就足以在任务执行时引发ModuleNotFoundError凭据与集群专属配置放入 connections / secrets backend不要在 DAG 源文件中硬编码认证信息使用场景有边界provider 面向 Airflow 编排 Hive 工作Airflow 之外的普通 Python 应用应使用专门的 Hive 客户端。总结apache-airflow-providers-apache-hive9.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点击查看免费下载相关推荐在 Apache Airflow DAG 中编排 Apache Beam 管道apache-airflow-providers-apache-beam 实战指南在 Apache Airflow DAG 中编排 Apache Beam 管道apache airflow providers apache beam 实战指在 Airflow DAG 中运行 Pig Latinapache-airflow-providers-apache-pig 4.8.2 实战指南在 Airflow DAG 中运行 Pig Latinapache airflow providers apache pig 4.8.2 实战指南 apachAirflow DAG 中提交 Spark 应用apache-airflow-providers-apache-spark 5.5.1 实战指南Airflow DAG 中提交 Spark 应用apache airflow providers apache spark 5.5.1 实战指南 本篇技术指南上一篇komorebi 的 container-padding 命令按工作区精确控制容器内边距下一篇brpc ExecutionQueue 深入指南wait-free 异步串行任务队列的设计与实战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?
咨询建站