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: ready、lifecycle: 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.1 | Bug Fixes | 使cncf.kubernetes成为 flink provider 的必需依赖(#29710) | — |
| 1.1.0 | Misc | 最低 Airflow 版本抬升至 2.4(#30917) | 2.4+ |
| 1.1.1 | Misc | 移除 Python 3.7 支持(#30963) | — |
| 1.1.2 | Misc | 简化 providers/apache 中len()条件写法(#33564) | — |
| 1.1.3 | Misc | 将部分模块 import 移入 type-checking 块以改善导入性能(#33754) | — |
| 1.2.0 | Bug Fixes + Misc | 修复soft_fail参数在异常抛出时未被尊重的问题(#34476);最低 Airflow 抬升至 2.5(#34728) | 2.5+ |
| 1.3.0 | Misc | 最低 Airflow 版本抬升至 2.6(#36017) | 2.6+ |
| 1.4.0 | Misc | 最低 Airflow 版本抬升至 2.7(#39240) | 2.7+ |
| 1.4.1 | Misc | 加速并简化airflow_version导入(#39497, #39552) | — |
| 1.4.2 | Misc | 实现按 Provider 使用最低直接依赖解析的测试(#39946) | — |
| 1.5.0 | Misc | 最低 Airflow 版本抬升至 2.8(#41396) | 2.8+ |
| 1.5.1 | Misc | 从 providers 中移除已废弃的soft_fail参数(#41710) | — |
| 1.6.0 | Misc | 最低 Airflow 版本抬升至 2.9(#44956);更新多 Provider 文档中的示例 DAG 链接 | 2.9+ |
| 1.6.1 | Misc | 将 flit 构建工具升级到 3.11.0(#46938);Apache Flink 迁移到新 Provider 目录结构(#46132) | — |
| 1.6.2 | Misc | 移除多余的else块(#49199) | — |
| 1.7.0 | Misc | 最低 Airflow 版本抬升至 2.10(#49843) | 2.10+ |
| 1.7.1 | Misc | 更新 BaseOperator 导入以兼容 Airflow 3.0(#52504);移除 Python 3.9 支持(#52072);改用 task-sdk 中的 BaseSensorOperator(#52296) | — |
| 1.7.2 | Bug Fixes + Misc | 修复FlinkKubernetesSensor._log_driver()中 pod 检索命名空间被硬编码的问题(#53653);新增 Python 3.13 支持(#46891) | — |
| 1.7.3 | Misc | Apache providers 迁移到common.compat(#57016) | — |
| 1.7.4 | Misc | 使所有 airflow 发行版符合 ASF 要求(#58138) | — |
| 1.8.0 | Misc | 最低 Airflow 版本抬升至 2.11(#58612) | 2.11+ |
| 1.8.1 | Misc | 为 providers 中的异常添加向后兼容支持(#58727) | — |
| 1.8.2 | Misc | 新年版权声明更新(#60344) | — |
| 1.8.3 | Misc | 将最低cryptography提升至 44.0.3、paramiko提升至 3.4.0(#62723) | — |
| 1.8.4 | Misc | 新增 Python 3.14 支持(#63520) | — |
| 1.8.5 | Misc | 为基于 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(任务成功);- 其余状态(
DEPLOYING、DEPLOYED_NOT_READY等)→ 返回False等待下次轮询。
对应地,测试文件中的test_cluster_error_state与test_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.3 | 2.4+ |
| 1.2.0 → 1.2.x | 2.5+ |
| 1.3.0 | 2.6+ |
| 1.4.0 → 1.4.2 | 2.7+ |
| 1.5.0 → 1.5.1 | 2.8+ |
| 1.6.0 → 1.6.2 | 2.9+ |
| 1.7.0 | 2.10+ |
| 1.7.1 → 1.7.4 | 2.10+ |
| 1.8.0 → 1.8.5 | 2.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 提升cryptography与paramiko最低版本(#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_name、namespace、attach_log(是否把 TaskManager Pod 日志追加到 Sensor 日志)、taskmanager_pods_namespace等之上,其轮询状态机与日志采集逻辑正是 1.7.2 修复的核心区域。
升级实战:如何基于 Changelog 制定升级策略
基于以上分析,升级 Flink Provider 时可以按以下清单逐项核对:
- 核对 Airflow 版本门槛:确认目标 Provider 版本的最低 Airflow 版本要求。若你的 Airflow 低于门槛,锁定旧 Provider 版本(如 Airflow 2.10 用户停留在 1.7.x);
- 核对 Python 版本:确认你的 Python 版本在目标 Provider 版本的支持窗口内(如 Python 3.14 需 1.8.4+,Python 3.9 最高只能用到 1.7.0);
- 检查已移除参数:若 DAG 中使用了
soft_fail参数,需在升级到 1.5.1+ 前移除; - 检查依赖升级:1.8.3 起要求
cryptography>=44.0.3、paramiko>=3.4.0,需评估环境中这两个库的兼容性; - 关注部署拓扑:若 Flink 应用部署在非默认命名空间且依赖日志采集,务必升级到 1.7.2+ 以获得命名空间修复;
- 核对行为语义:阅读 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),仅供参考