Apache Airflow Flink Provider 版本演进全解:从 1.0.0 到 1.8.5 的 Changelog 深度剖析
2026/9/13 2:51:09 网站建设 项目流程

Apache Airflow Flink Provider 版本演进全解:从 1.0.0 到 1.8.5 的 Changelog 深度剖析

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

apache-airflow-providers-apache-flink是 Apache Airflow 官方维护的社区托管 Provider,用于在 Kubernetes 集群上通过 Flink Kubernetes Operator 以自定义资源(FlinkDeployment)方式程序化提交、调度与监控 Apache Flink 流式作业。本文以 providers/apache/flink/docs/changelog.rst 为骨架,逐版本梳理从 1.0.0 首发到 1.8.5 的全部变更记录,并结合仓库内 Operator/Sensor 源码与单元测试,讲清每个关键修复的底层原理、每次最低 Airflow 版本抬升背后的支持政策,以及 Python 版本支持范围的迁移轨迹,帮助你在升级 Provider 时做出有依据的决策。

Changelog 文档的定位与阅读方式

Apache Airflow 每个社区托管 Provider 都维护一份独立的 Changelog,存放在providers/<厂商>/<产品>/docs/changelog.rst。这份文档由 release manager 半自动维护,因此其内容约定非常明确:

  • 文档头部注释(providers/apache/flink/docs/changelog.rst)说明:只有当存在 breaking changes、需要向用户解释应对方式时,才在 "Changelog" 标题下方直接添加说明条目;
  • 每个版本号下面按变更类型分组,Flink Provider 主要出现两种分类:Bug Fixes(缺陷修复)与Misc(杂项、依赖升级、工程化改造);
  • .. Below changes are excluded from the changelog.注释标记的行是被排除的变更——它们通常是纯内部流程(如"准备文档"、"移除多余 LICENSE 文件"、"切换 pre-commit 工具"),对用户无感知,不进入正式变更清单。这份文档本身就是阅读 Provider 演变史的第一手资料。

从 provider.yaml 可以看到,当前 Provider 状态为state: readylifecycle: production,即已进入生产可用阶段;README.rst 表明最新发布版本为1.8.5

快速上手:包安装与运行环境要求

Flink Provider 不是 Airflow 核心自带的模块,需要独立安装。根据 README.rst 与 docs/index.rst,安装方式为:

pip install apache-airflow-providers-apache-flink

该包支持在现有 Airflow 安装之上叠加安装,当前版本(1.8.5)的依赖要求如下:

PIP 包最低版本要求
apache-airflow>=2.11.0
apache-airflow-providers-common-compat>=1.10.1
cryptography>=44.0.3
apache-airflow-providers-cncf-kubernetes>=5.1.0

同时支持 Python 版本:3.10、3.11、3.12、3.13、3.14

注意:这些要求是当前最新版本的门槛,历史版本的依赖要求随版本演进不断变化(详见下文"最低 Airflow 版本支持政策"小节)。安装时 pip 会自动解析出与你的 Airflow 版本兼容的 Provider 版本。

版本演进全景:27 个版本、26 次发布

结合 provider.yaml 中维护的版本清单与 changelog,Flink Provider 自 1.0.0 起共发布了 27 个版本。下表将每个版本的核心变更、最低 Airflow 版本与 Python 支持变化汇总为一张速查表:

版本变更类型核心内容最低 Airflow 版本
1.0.0首发Initial version of the Flink provider
1.0.1Bug Fixes使cncf.kubernetes成为 flink provider 的必需依赖(#29710)
1.1.0Misc最低 Airflow 版本抬升至 2.4(#30917)2.4+
1.1.1Misc移除 Python 3.7 支持(#30963)
1.1.2Misc简化 providers/apache 中len()条件写法(#33564)
1.1.3Misc将部分模块 import 移入 type-checking 块以改善导入性能(#33754)
1.2.0Bug Fixes + Misc修复soft_fail参数在异常抛出时未被尊重的问题(#34476);最低 Airflow 抬升至 2.5(#34728)2.5+
1.3.0Misc最低 Airflow 版本抬升至 2.6(#36017)2.6+
1.4.0Misc最低 Airflow 版本抬升至 2.7(#39240)2.7+
1.4.1Misc加速并简化airflow_version导入(#39497, #39552)
1.4.2Misc实现按 Provider 使用最低直接依赖解析的测试(#39946)
1.5.0Misc最低 Airflow 版本抬升至 2.8(#41396)2.8+
1.5.1Misc从 providers 中移除已废弃的soft_fail参数(#41710)
1.6.0Misc最低 Airflow 版本抬升至 2.9(#44956);更新多 Provider 文档中的示例 DAG 链接2.9+
1.6.1Misc将 flit 构建工具升级到 3.11.0(#46938);Apache Flink 迁移到新 Provider 目录结构(#46132)
1.6.2Misc移除多余的else块(#49199)
1.7.0Misc最低 Airflow 版本抬升至 2.10(#49843)2.10+
1.7.1Misc更新 BaseOperator 导入以兼容 Airflow 3.0(#52504);移除 Python 3.9 支持(#52072);改用 task-sdk 中的 BaseSensorOperator(#52296)
1.7.2Bug Fixes + Misc修复FlinkKubernetesSensor._log_driver()中 pod 检索命名空间被硬编码的问题(#53653);新增 Python 3.13 支持(#46891)
1.7.3MiscApache providers 迁移到common.compat(#57016)
1.7.4Misc使所有 airflow 发行版符合 ASF 要求(#58138)
1.8.0Misc最低 Airflow 版本抬升至 2.11(#58612)2.11+
1.8.1Misc为 providers 中的异常添加向后兼容支持(#58727)
1.8.2Misc新年版权声明更新(#60344)
1.8.3Misc将最低cryptography提升至 44.0.3、paramiko提升至 3.4.0(#62723)
1.8.4Misc新增 Python 3.14 支持(#63520)
1.8.5Misc为基于 flit 的 pyproject.toml 添加显式[tool.flit.sdist]段(#65861)

从这张表可以清晰地看出两条主线:功能与兼容性持续增强(Python 支持从 3.7 一路扩展到 3.14,Airflow 最低版本从 2.4 抬升到 2.11),以及工程化与合规性持续完善(ASF 合规、构建工具、测试基础设施)。

关键 Bug Fix 源码级解析

1.7.2:_log_driver()命名空间硬编码修复

这是 changelog 中最重要的缺陷修复(changelog.rst):fix hardcoded value of FlinkKubernetesSensor._log_driver() namespace for pod retrieval (#53653)

问题背景FlinkKubernetesSensor支持attach_log=True将 TaskManager Pod 的日志追加到 Sensor 日志中。_log_driver()负责按 TaskManager 的labelSelector列出 Pod 并拉取日志。修复前,Pod 检索的命名空间是硬编码值;当 FlinkDeployment 部署在非默认命名空间时,get_namespaced_pod_list会在错误的命名空间中查询 Pod,导致日志拉取失败。

修复后的源码行为(sensors/flink_kubernetes.py):

all_pods = self.hook.get_namespaced_pod_list( namespace=self.taskmanager_pods_namespace or self.namespace or "default", watch=False, label_selector=task_manager_labels, )

命名空间解析优先级为:taskmanager_pods_namespace(显式指定的 TaskManager Pod 命名空间)→namespace(FlinkDeployment 所在命名空间)→"default"。与此同时,日志读取目标命名空间取自响应中的response["metadata"]["namespace"](L97),即部署实际所在的命名空间。

测试佐证:单元测试文件 tests/unit/apache/flink/sensors/test_flink_kubernetes.py 提供了多组针对命名空间行为的用例,包括:

  • test_namespace_from_sensor:从 Sensor 参数取命名空间;
  • test_namespace_from_connection:从 Kubernetes Connection 的extra__kubernetes__namespace取命名空间(该测试还构造了kubernetes_with_namespace连接,extra 中携带"mock_namespace");
  • test_namespace_from_taskmanager相关用例:验证taskmanager_pods_namespace=namespae_name参数传递路径(L1154-L1162)。

这些用例共同保障了"跨命名空间部署 Flink 应用并采集日志"这一典型生产场景的正确性。

1.2.0:尊重soft_fail参数

另一个值得注意的修复出现在 1.2.0(changelog.rst):fix(providers/flink): respect soft_fail argument when exception is raised (#34476)

问题本质soft_fail=True是 Airflow BaseSensorOperator 提供的机制——当 Sensor 判定任务失败时,不抛出硬异常导致任务失败,而是让任务以skipped状态优雅跳过。修复前,Flink 传感器在检测到失败状态时直接抛异常,未走 soft_fail 分支。

当前实现(sensors/flink_kubernetes.py)中,poke()方法按jobManagerDeploymentStatus状态机推进:

FAILURE_STATES = ("MISSING", "ERROR") SUCCESS_STATES = ("READY",) def poke(self, context: Context) -> bool: response = self.hook.get_custom_object(...) application_state = response["status"]["jobManagerDeploymentStatus"] ... if application_state in self.FAILURE_STATES: message = f"Flink application failed with state: {application_state}" raise AirflowException(message) if application_state in self.SUCCESS_STATES: self.log.info("Flink application ended successfully") return True self.log.info("Flink application is still in state: %s", application_state) return False
  • 状态缺失(KeyError)→ 返回False(继续轮询);
  • MISSING/ERROR→ 抛出AirflowException(由 BaseSensorOperator 依据soft_fail决定是失败还是跳过);
  • READY→ 返回True(任务成功);
  • 其余状态(DEPLOYINGDEPLOYED_NOT_READY等)→ 返回False等待下次轮询。

对应地,测试文件中的test_cluster_error_statetest_missing_cluster都断言pytest.raises(AirflowException)(L910-L923、L985-L998)。

演进轨迹soft_fail在 1.2.0 被正确支持后,于 1.5.1 被统一移除(#41710)——这是 Apache Airflow 全局性的 API 清理:soft_fail参数从所有 providers 中删除(该参数的行为由更通用的重试/跳过机制替代)。如果你的 DAG 还在使用soft_fail=True,升级到 1.5.1 及以上版本时需要注意该参数已不存在。

最低 Airflow 版本支持政策:为什么每个大版本都要"抬门槛"

changelog 中反复出现一类条目:"Bump minimum Airflow version in providers to Airflow 2.x"。这是 Apache Airflow 社区对社区托管 Provider 的统一支持政策:随着 Airflow 核心版本迭代,旧版 Airflow 不再获得安全修复与兼容性保障,Provider 的最低支持版本会周期性上移。

从 changelog 可以还原出 Flink Provider 的最低 Airflow 版本演进轨迹:

Provider 版本最低 Airflow 版本
1.0.x – 1.0.1未声明(随 Airflow 2.2+ 生态)
1.1.0 → 1.1.32.4+
1.2.0 → 1.2.x2.5+
1.3.02.6+
1.4.0 → 1.4.22.7+
1.5.0 → 1.5.12.8+
1.6.0 → 1.6.22.9+
1.7.02.10+
1.7.1 → 1.7.42.10+
1.8.0 → 1.8.52.11+

以 1.8.0 为例,changelog 特别以 note 形式强调:该版本仅适用于 Airflow 2.11+(changelog.rst)。这意味着如果你仍在使用 Airflow 2.10 及以下版本,升级 Flink Provider 到 1.8.0 会导致 pip 依赖解析失败,应当锁定在 1.7.x 系列。

源码层面的兼容机制:Provider 包内维护了 version_compat.py,它通过packaging.version解析airflow.__version__,并导出AIRFLOW_V_3_0_PLUS常量,供代码在 Airflow 2.x 与 3.x 之间做条件分支。这与 1.7.1 中"Update BaseOperator imports for Airflow 3.0 compatibility"(#52504)、"Use BaseSensorOperator from task sdk in providers"(#52296)等变更互为印证——Provider 代码通过airflow.providers.common.compat.sdk统一导入BaseOperator/BaseSensorOperator,屏蔽了 Airflow 2.x 与 3.x 在算子基类实现上的差异。

Python 版本支持范围迁移

changelog 同时记录了 Flink Provider 对 Python 版本支持范围的完整变化:

  • 1.1.1(#30963):移除 Python 3.7 支持;
  • 1.7.1(#52072):移除 Python 3.9 支持;
  • 1.7.2(#46891):新增 Python 3.13 支持;
  • 1.8.4(#63520):新增 Python 3.14 支持。

结合当前 README.rst 声明的支持范围3.10 – 3.14,可以推断出该 Provider 的 Python 支持窗口大致保持在"最新 4 个稳定版本"。升级 Python 版本前,请核对你的 Provider 版本是否在支持窗口内。

工程化与生态演进:从打包、依赖到合规

除去功能与兼容性,changelog 中大量Misc条目还勾勒出 Apache Airflow 整个 Provider 生态的工程化演进,Flink Provider 是这一进程的参与者:

  • 打包工具链:1.6.1 升级 flit 至 3.11.0(#46938),1.8.5 为基于 flit 的 pyproject.toml 添加显式[tool.flit.sdist]段(#65861),保证 sdist 构建的可复现性;
  • 项目结构迁移:1.6.1 "Move Apache Flink to new provider structure"(#46132),随后 1.6.2 移除多余else块(#49199)等代码清理;
  • 依赖治理:1.0.1 将apache-airflow-providers-cncf-kubernetes从可选依赖提升为必需依赖(#29710)——这是符合直觉的设计,因为 Flink Provider 的所有 Operator/Sensor 都构建在KubernetesHook之上(见 operators/flink_kubernetes.py 与 sensors/flink_kubernetes.py 的 import);1.8.3 提升cryptographyparamiko最低版本(#62723)则是对安全漏洞的响应式升级;
  • 测试基础设施:1.4.2 实现按 Provider 的最低直接依赖解析测试(#39946),确保每个 Provider 在最小依赖集下也能通过测试;
  • ASF 合规:1.7.4 "Convert all airflow distributions to be compliant with ASF requirements"(#58138)、1.8.2 更新版权声明(#60344),属于 Apache 基金会层面的规范要求;
  • 运行期性能:1.1.3 将部分 import 移入TYPE_CHECKING块(#33754)、1.4.1 简化airflow_version导入(#39497, #39552),都是减少 Provider 加载开销的优化。

组件速览:Changelog 背后的两个核心类

理解 changelog 中 Bug Fix 的意义,需要先知道 Flink Provider 的全部对外能力。该包只有两个核心组件(见 provider.yaml):

FlinkKubernetesOperator(operators/flink_kubernetes.py):在 Kubernetes 集群中创建flinkDeployment自定义资源对象。核心参数包括:

  • application_file:FlinkDeployment 的 CRD 定义,支持.yaml/.json文件路径或 YAML/JSON 字符串(可模板化,template_fields = ("application_file", "namespace"));
  • namespace:部署 FlinkDeployment 的 Kubernetes 命名空间;
  • kubernetes_conn_id:Kubernetes 连接 ID,默认kubernetes_default
  • api_group/api_version:CRD 的 API 组与版本,默认flink.apache.org/v1beta1
  • in_cluster/cluster_context/config_file:Kubernetes 客户端认证配置;
  • plural:自定义资源复数名,默认flinkdeployments

execute()流程(L102-L118)先通过list_cluster_custom_object校验 API 可达性,再调用create_custom_object提交 CRD 并返回响应。

FlinkKubernetesSensor(sensors/flink_kubernetes.py):轮询 FlinkDeployment 的状态直至READY或失败。核心参数在application_namenamespaceattach_log(是否把 TaskManager Pod 日志追加到 Sensor 日志)、taskmanager_pods_namespace等之上,其轮询状态机与日志采集逻辑正是 1.7.2 修复的核心区域。

升级实战:如何基于 Changelog 制定升级策略

基于以上分析,升级 Flink Provider 时可以按以下清单逐项核对:

  1. 核对 Airflow 版本门槛:确认目标 Provider 版本的最低 Airflow 版本要求。若你的 Airflow 低于门槛,锁定旧 Provider 版本(如 Airflow 2.10 用户停留在 1.7.x);
  2. 核对 Python 版本:确认你的 Python 版本在目标 Provider 版本的支持窗口内(如 Python 3.14 需 1.8.4+,Python 3.9 最高只能用到 1.7.0);
  3. 检查已移除参数:若 DAG 中使用了soft_fail参数,需在升级到 1.5.1+ 前移除;
  4. 检查依赖升级:1.8.3 起要求cryptography>=44.0.3paramiko>=3.4.0,需评估环境中这两个库的兼容性;
  5. 关注部署拓扑:若 Flink 应用部署在非默认命名空间且依赖日志采集,务必升级到 1.7.2+ 以获得命名空间修复;
  6. 核对行为语义:阅读 docs/operators.rst 中的FlinkKubernetesOperator使用指南,确认参数语义未变化。

总结

apache-airflow-providers-apache-flink的 Changelog 完整记录了该 Provider 从 2022 年首发(1.0.0)到 1.8.5 的三年演进:功能上以 Flink Kubernetes Operator 自定义资源为支点,通过FlinkKubernetesOperator提交、FlinkKubernetesSensor轮询的方式支撑流式作业的生命周期管理;工程上则紧随 Airflow 生态完成了最低版本抬升、Python 窗口迁移、task-sdk 兼容、ASF 合规等系统性改造。对于使用者而言,这份 changelog 既是升级决策的依据,也是理解 Provider 内部实现演变的最佳索引——结合 源码 与 单元测试 阅读,可以准确掌握每个版本变更背后的真实行为差异。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询