尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Airflow Apache HDFS Provider 版本演进与核心能力全解析:从 snakebite 迁移到 WebHDFS 与远程日志体系

Apache Airflow Apache HDFS Provider 版本演进与核心能力全解析:从 snakebite 迁移到 WebHDFS 与远程日志体系 Apache Airflow Apache HDFS Provider 版本演进与核心能力全解析从 snakebite 迁移到 WebHDFS 与远程日志体系【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的apache-airflow-providers-apache-hdfs是官方维护的 HDFS 集成 Provider负责在 Airflow 工作流中对接 Hadoop 分布式文件系统。本文以该 Provider 的官方变更日志changelog.rst为骨架完整梳理其从 1.0.0 到 4.13.0 的演进历程并结合当前仓库源码深入剖析WebHDFSHook、WebHDFS 传感器、HDFS 远程日志写入、HDFS Asset 资产管理等核心能力帮助你在选型、升级与日常使用中做出正确决策。一、Provider 概览与安装前提apache-airflow-providers-apache-hdfs是 Apache Airflow 社区管理的 Provider 之一所有类均位于airflow.providers.apache.hdfsPython 包内当前发布版本为4.13.0见 index.rst。在现有 Airflow 安装之上通过 pip 即可安装pip install apache-airflow-providers-apache-hdfs根据 index.rst 中记录的依赖矩阵4.13.0 对运行环境有以下要求依赖包版本要求apache-airflow2.11.0apache-airflow-providers-common-compat1.12.0hdfs[avro,dataframe,kerberos]2.5.4Python 3.122.7.3Python 3.12fastavro1.10.0Python 3.131.12.1Python 3.14pandas2.1.2Python 3.132.2.3Python 3.132.3.3Python 3.14依赖hdfs即 hdfscli 客户端库它是 WebHDFS 协议实现的基石。从 4.11.0 起Provider 的最低 Airflow 版本被提升到2.11.0此前依次为 2.104.9.0、2.94.7.0、2.84.5.0、2.74.4.0、2.64.3.0、2.54.2.0、2.44.0.0等遵循 Apache Airflow 社区 Provider 的最低版本支持策略。二、版本演进主线从变更日志看能力变迁2.1 4.0.0 里程碑告别 snakebite全面转向 WebHDFS变更日志中最重要的转折点是4.0.0的 Breaking changes基于旧snakebite-py3库的 HDFS Hook 与 Sensor 被彻底移除包括HDFSHook、HDFSSensor、HdfsRegexSensor等变更 PR #31262。移除理由是 snakebite-py3 多年未更新且其依赖的 protobuf 3 已在 2023 年 6 月终止生命周期与当时主流环境中的 protobuf 4 不兼容。对于仍在使用 3.* 版本的用户变更日志给出了两条临时性 workaround设置环境变量PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATIONpython或将 protobuf 降级到最新的 3.* 版本当时为 3.20.3。但二者均有明显代价前者会让大量使用 protobuf 的库含多个 Google 客户端库与 Kubernetes运行显著变慢后者会与新版 Google、gRPC 等 Provider 产生冲突。因此变更日志明确建议这只是临时方案应尽快迁移到基于 WebHDFS API 的WebHDFSHook与WebHDFSSensor。这条主线决定了 Provider 4.x 时代的全部技术方向。2.2 远程日志体系从 4.1.0 到 4.13.0 的持续演进HDFS 远程日志是 Provider 近几个版本投入最大的方向4.1.0新增从 HDFS 读写任务实例日志的能力PR #31512首次打通HdfsTaskHandler4.5.1移除已废弃的日志处理器参数filename_templatePR #41552接口进一步收敛4.8.0正式支持在Airflow 3中把任务日志读写到 HDFSPR #48788并同步适配了新架构4.10.3HdfsTaskHandler与其他众多远程日志 Handler 一起支持日志文件大小处理max_bytes/backup_countPR #554554.12.1修复远程日志 Provider 不满足RemoteLogIO上传契约的问题PR #683004.13.0新增HdfsRemoteLogIO.from_config并注册hdfs远程日志 schemePR #71278使 HDFS 远程日志接入方式与 S3/GCS 等其他 Provider 对齐。启用 HDFS 远程日志的配置方式摘自 logging/hdfs-task-handler.rst在airflow.cfg中如下[logging] # Airflow can store logs remotely in HDFS. Users must supply a remote # location URL (starting with either hdfs://...) and an Airflow connection # id that provides access to the storage location. remote_logging True remote_base_log_folder hdfs://some/path/to/logs remote_log_conn_id webhdfs_default远程日志读写依赖一个已正确配置的 Airflow 连接连接未就绪时整个过程会失败。示例中 Airflow 将尝试使用WebHDFSHook(webhdfs_default)。2.3 连接与安全能力增强WebHDFS 连接在多个版本中持续强化2.1.0支持 SSL WebHDFS 连接PR #176372.2.0恢复 WebHDFS 的 HA高可用支持PR #197113.1.0为 WebHDFS Sensor 增加认证能力PR #251103.0.1修复WebHDFSHook的可选端口问题PR #245504.3.1修复webhdfs连接类型注册PR #361454.7.0为WebHDFSHook增加mTLS支持PR #445614.10.0在WebHDFSHook中支持自定义 headers 与 cookiesPR #50955并改用conn.password读取密码PR #50759。2.4 传感器与资产管理4.6.0新增MultipleFilesWebHdfsSensorPR #43045可等待目录中的多个文件全部就位4.12.0为新 scheme 增加 URI 清洗器sanitizer与 Asset 工厂PR #66426并修复 Asset 归一化的若干问题PR #66710——这对应 Airflow 3 的 Asset原 Dataset体系4.13.0注册hdfs远程日志 scheme 的同时Asset 相关能力也随 Provider 信息一并发布。2.5 维护性变更Misc与版本节奏变更日志中大量的 Misc 条目揭示了 Provider 的工程化演进最低 Python 版本随上游策略不断上调如 4.10.1 移除 Python 3.9 支持、4.11.4 增加 Python 3.14 支持、4.10.2 增加 Python 3.13 支持BaseHook、BaseOperator、BaseSensorOperator逐步迁移到 Task SDKPR #51873、#52505、#52296Provider 元数据改为从 YAML 加载PR #63826避免导入 Hook 类即可读取 Hook 信息。这些变更对使用者基本透明但反映了 Provider 与 Airflow 3 新架构对齐的整体趋势。三、源码级解读4.13.0 的核心实现3.1 WebHDFSHook连接、HA 与安全协议WebHDFSHook位于 hooks/webhdfs.py是对 hdfscli 客户端的封装。其连接参数包括webhdfs_conn_id默认webhdfs_default与proxy_user代理认证用户连接类型为webhdfs。Hook 的 HA 逻辑体现在_find_valid_server方法中连接 Host 支持逗号分隔的多个 NameNode 地址Hook 会依次尝试 TCP 建连socket.connect_ex并对每个可达节点执行client.status(/)探活首个校验成功的 NameNode 即被选中失败节点会被记录警告并跳过。_get_client方法根据安全模式构建客户端Kerberos 模式当core/security配置为kerberos时优先导入并使用hdfs.ext.kerberos.KerberosClient若扩展缺失则报错退出非 Kerberos 模式使用InsecureClient认证用户取proxy_user或连接登录名。连接extra字段支持以下可选参数同时记录在 connections.rst 中参数含义use_ssl是否启用 SSL默认false启用后连接串切换为https://verifySSL 证书校验方式CA 证书路径或布尔值默认False对应 requests 的verifycertmTLS 客户端证书路径可单独使用或与key组合key与cert搭配的 mTLS 客户端私钥路径cookies附加到请求会话的 cookiesheaders附加到请求会话的 headersHook 对外暴露三个主要操作方法check_for_path(hdfs_path)通过conn.status(path, strictFalse)查询 FileStatus返回路径是否存在这是传感器的核心依赖load_file(source, destination, overwriteTrue, parallelism1)上传本地文件或文件夹到 HDFSparallelism控制并行线程数0或负数表示按文件数自动开线程若源是文件夹仅上传其中非空文件read_file(filename)读取 HDFS 文件内容并返回原始字节串。3.2 WebHdfsSensor 与 MultipleFilesWebHdfsSensor两个传感器都定义在 sensors/web_hdfs.pyWebHdfsSensor通过filepath参数指定目标路径poke()时调用WebHDFSHook.check_for_path轮询直到文件或文件夹出现在 HDFS 中filepath是模板字段支持 Jinja 渲染MultipleFilesWebHdfsSensor4.6.0 引入接收directory_path与expected_filenamespoke()中通过conn.list(directory_path)列出目录实际文件与期望文件集合做差集存在缺失文件则返回False继续等待。directory_path与expected_filenames均为模板字段。两个传感器均继承自 Task SDK 的BaseSensorOperator可直接用于 Airflow 3 的任务装饰器或传统 DAG 写法。3.3 HdfsTaskHandler 与 HdfsRemoteLogIO远程日志读写链路远程日志实现位于 log/hdfs_task_handler.py由两个核心类组成HdfsRemoteLogIO是 Airflow 3 的RemoteLogIO接口实现负责与 HDFS 的实际读写from_config()4.13.0 新增从 Airflow 日志配置构建实例——读取logging/remote_task_handler_kwargsJSON 配置过滤出属于 IO 层的参数同时读取base_log_folder、remote_base_log_folder取 URL 的 path 部分与delete_local_logsupload()将本地日志文件上传到 HDFS 的remote_base对应路径上传成功后若开启delete_local_copy则删除本地副本read()按remote_base 相对路径在 HDFS 中查找日志存在则读取并解码为 UTF-8 文本否则返回No logs found on hdfs for ti...提示信息内部通过WebHDFSHook连接 ID 来自logging/REMOTE_LOG_CONN_ID配置完成实际 I/O。HdfsTaskHandler继承 Airflow 的FileTaskHandler是注册到日志系统的处理器支持max_bytes与backup_count参数即 4.10.3 引入的日志文件大小处理能力set_context()计算本地日志相对路径并在任务重跑如传感器重调度时先清空本地文件避免重复数据上传close()时通过io.upload()将日志上传到 HDFS并通过closed标志防止logging.shutdown触发重复上传_read_remote_logs()在读取时按任务实例与重试序号渲染远程日志相对路径从 HDFS 拉取日志供 UI 展示。在 provider.yaml 中该能力注册为logging: - airflow.providers.apache.hdfs.log.hdfs_task_handler.HdfsTaskHandler remote-logging: - classpath: airflow.providers.apache.hdfs.log.hdfs_task_handler.HdfsRemoteLogIO scheme: hdfshdfsscheme 的注册4.13.0意味着 Airflow 3 可自动将hdfs://开头的remote_base_log_folder路由到HdfsRemoteLogIO。3.4 HDFS AssetAirflow 3 数据资产管理assets/hdfs.py 实现了 Airflow 3 的 Asset即原 Dataset支持sanitize_uri()校验hdfs://URI 必须包含 NameNode 主机与路径否则抛出ValueErrorcreate_asset()按host、path、port默认 8020构造hdfs://host:port/path形式的 Asset URIconvert_asset_to_openlineage()将 Asset URI 转换为 OpenLineage Datasetnamespace 为hdfs://netlocname 为去首斜杠的路径。以上三项在 provider.yaml 中分别以asset-urisAirflow 3与dataset-uris向后兼容两种 scheme 注册handler/factory/to_openlineage_converter 一一对应与 4.12.0 变更日志中的「URI sanitizer 与 Asset factory」吻合。四、升级与兼容性实操指南4.1 从 3.* 迁移到 4.x 的检查清单替换 Hook/Sensor 类将所有HDFSHook、HDFSSensor、HdfsRegexSensor引用替换为WebHDFSHook、WebHdfsSensor重配连接连接类型改为webhdfsHost 支持逗号分隔的多 NameNode 列表HA 场景按需补充use_ssl/verify/cert/key/cookies/headers等 extra 参数确认 Airflow 版本4.13.0 要求 Airflow 2.11.0升级前请核对评估认证方式Kerberos 集群需保证hdfs[kerberos]扩展可用非 Kerberos 环境通过 Login有效用户与proxy_user控制执行身份。4.2 版本能力速查表版本类型关键内容4.13.0FeatureHdfsRemoteLogIO.from_config注册hdfs远程日志 scheme4.12.1Bug Fix修复 RemoteLogIO 上传契约问题4.12.0Feature/Bug FixURI sanitizer 与 Asset factoryAsset 归一化修复4.10.3Bug Fix多远程日志 Handler 支持日志文件大小处理4.10.0FeatureWebHDFSHook 支持自定义 headers/cookies4.8.0FeatureAirflow 3 中读写任务日志到 HDFS4.7.0FeatureWebHDFSHook 增加 mTLS 支持4.6.0Feature新增MultipleFilesWebHdfsSensor4.1.0Feature读写任务实例日志到 HDFS首个版本4.0.0Breaking移除 snakebite-py3 的 HDFS Hook/Sensor迁移到 WebHDFS五、总结从变更日志可以清晰看到apache-airflow-providers-apache-hdfs的技术收敛路径4.0.0 果断废弃长期失修的 snakebite 方案将全部能力收敛到 WebHDFS 协议之上随后围绕 WebHDFS 持续增强安全SSL、mTLS、headers/cookies与高可用多 NameNode 探活在 Airflow 3 时代则同步完成了远程日志读写、RemoteLogIO 契约对齐、hdfsscheme 注册以及 Asset 资产体系的接入。理解这条演进主线有助于你在选型时直接采用 WebHDFS 方案、在升级时规避 breaking changes并快速上手 HDFS 远程日志与资产管理等新特性。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表