ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Airflow common.compat 兼容 Provider 详解:跨 Airflow 2.x/3.x 的兼容层设计与使用

Apache Airflow common.compat 兼容 Provider 详解:跨 Airflow 2.x/3.x 的兼容层设计与使用 Apache Airflow common.compat 兼容 Provider 详解跨 Airflow 2.x/3.x 的兼容层设计与使用【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowapache-airflow-providers-common-compat简称common.compat是 Apache Airflow 官方 Provider 体系中的通用兼容包它的职责是提供一段段可在旧版与新版 Airflow 之间复用的兼容性代码。阅读本文后你将了解该包的定位、安装方式与版本要求、可选依赖extras配置并能深入其源码掌握基于 PEP 562 懒加载导入的__getattr__回退机制、version_compat版本常量、Dataset→Asset 重命名兼容以及异步 Connection/Hook 的跨版本实现从而为自己的 Operator、Hook 或 Plugin 编写同时兼容 Airflow 2.x 与 3.x 的代码。包定位为旧版本 Airflow 提供兼容性代码根据包索引文档 providers/common/compat/docs/index.rst该分发包的定义为Common Compatibility Provider - providing compatibility code for previous Airflow versions通用兼容 Provider——为先前版本的 Airflow 提供兼容性代码包名apache-airflow-providers-common-compat当前仓库中版本为1.19.0所有类都位于 Python 包airflow.providers.common.compat之下Provider 元数据状态、生命周期、历史版本列表维护在 provider.yaml其中state: ready、lifecycle: production并从 1.0.0 一路发布到 1.19.0构建配置见 pyproject.toml构建后端为flit_core4.0.2requires-python 3.10支持 Python 3.103.14见classifiers。这个包与一般 Provider 不同它不对接任何外部系统而是服务于 Airflow 生态本身——尤其是 Airflow 2 到 Airflow 3 的大版本迁移期间DAG 作者和第三方 Provider 开发者用它来避免在自己的代码中散落大量try/except ImportError。安装与版本要求在已有 Airflow 安装之上直接通过 pip 安装即可与 index.rst 中 Installation 章节一致pip install apache-airflow-providers-common-compat最低支持版本本 Provider 分发包支持的最低 Apache Airflow 版本为2.11.0。依赖要求同时与 pyproject.toml 中dependencies字段互相印证PIP 包版本要求apache-airflow2.11.0asgiref2.3.0Python 3.14asgiref3.11.1Python 3.14asgiref之所以是硬依赖是因为兼容层中的异步辅助函数如get_async_connection在其回退路径上使用了asgiref.sync.sync_to_async把同步实现包成协程详见后文。可选的跨 Provider 依赖cross-provider extras某些功能需要额外安装对应 Provider 分发才能启用可通过 extras 形式安装例如pip install apache-airflow-providers-common-compat[openlineage]依赖包Extraapache-airflow-providers-openlineageopenlineageapache-airflow-providers-standardstandard这一配置同样可见于 pyproject.toml 的[project.optional-dependencies]小节。之所以拆出standard这个 extra是因为common.compat中的部分兼容模块如standard/operators.py需要在 Airflow 2.x 上回退到airflow.providers.standard或airflow.operators.python的旧路径。官方发行包下载与校验官方发布的 sdist 与 wheel 包含.asc签名与.sha512校验和可从 Apache 官方下载站点获取以 1.19.0 为例apache_airflow_providers_common_compat-1.19.0.tar.gz与apache_airflow_providers_common_compat-1.19.0-py3-none-any.whl。从源码构建安装的通用流程可参考 从源码安装 Provider 指南。核心机制一version_compat 版本常量version_compat.py 是兼容层的地基。文件头部有一段重要注释说明了它的存在理由此文件被刻意手工复制到其他 Provider 中以避免在 Provider 之间引入不必要的依赖。如果你想在 Provider 中添加依赖 Airflow 版本的条件代码请把该文件复制到你 Provider 的根包中并从中导入这些常量。其实现非常直白——读取运行时 Airflow 的__version__并解析为三元组然后导出四个布尔常量def get_base_airflow_version_tuple() - tuple[int, int, int]: from packaging.version import Version from airflow import __version__ airflow_version Version(__version__) return airflow_version.major, airflow_version.minor, airflow_version.micro AIRFLOW_V_3_0_PLUS: bool get_base_airflow_version_tuple() (3, 0, 0) AIRFLOW_V_3_1_PLUS: bool get_base_airflow_version_tuple() (3, 1, 0) AIRFLOW_V_3_2_PLUS: bool get_base_airflow_version_tuple() (3, 2, 0) AIRFLOW_V_3_3_PLUS: bool get_base_airflow_version_tuple() (3, 3, 0)这些常量在包内被多处使用例如 standard/operators.py 用AIRFLOW_V_3_2_PLUS决定是直接导入 3.2 才存在的BaseAsyncOperator/is_async_callable还是在低版本上提供本地桩实现assets/init.py 用AIRFLOW_V_3_0_PLUS决定从airflow.datasets还是airflow.sdk.definitions.asset导入。文件注释中还提到一个历史细节BaseOperator曾被移出version_compat以避免循环导入。核心机制二基于 PEP 562 的懒加载导入回退sdk.py 是整个兼容层的核心模块其 docstring 点明了设计目标提供先尝试 Airflow 3 路径、再回退 Airflow 2 路径的懒导入使同一份代码在两个大版本上都能工作。三类映射表sdk.py内维护了三张表最终交给共享工具函数_compat_utils.create_module_getattr定义于 _compat_utils.py生成模块级__getattr___IMPORT_MAP普通导入映射类名 - 模块路径路径可以是单个字符串也可以是从新到旧依次尝试的元组。例如BaseOperator: (airflow.sdk, airflow.models.baseoperator), XCom: (airflow.sdk.execution_time.xcom, airflow.models.xcom), chain: (airflow.sdk, airflow.models.baseoperator), redact: ( airflow.sdk.log, airflow.sdk._shared.secrets_masker, airflow.sdk.execution_time.secrets_masker, airflow.utils.log.secrets_masker, ),覆盖范围包括 HooksBaseHook、OperatorBaseOperator、装饰器task/dag/task_group/setup/teardown、模型DAG/Connection/Variable/Param/XComArg、状态枚举DagRunState/TaskInstanceState/TriggerRule/WeightRule、血缘HookLineage*、异常AirflowException家族、TaskDeferred、XComNotFound等、可观测性Stats、监听器hookimpl、get_listener_manager以及配置conf等大类。_RENAME_MAP处理改名场景格式为新名字 - (新路径, 旧路径, 旧名字)。最典型的例子是 Airflow 3.0 中 Dataset 家族整体改名为 AssetAsset: (airflow.sdk, airflow.datasets, Dataset), AssetAlias: (airflow.sdk, airflow.datasets, DatasetAlias), AssetAll: (airflow.sdk, airflow.datasets, DatasetAll), AssetAny: (airflow.sdk, airflow.datasets, DatasetAny),_MODULE_MAP整模块级回退例如timezoneairflow.sdk.timezone回退到airflow.utils.timezone与ioairflow.sdk.io回退到airflow.io使用方式为from airflow.providers.common.compat.sdk import timezone。此外还有一个_AIRFLOW_3_ONLY_EXCEPTIONS字典DownstreamTasksSkipped、DagRunTriggerException这类只在 Airflow 3 中存在的异常会在运行时检测到AIRFLOW_V_3_0_PLUS为真时才动态并入_IMPORT_MAPAssetLineageInfo的模块位置在 3.0/3.1 与 3.2 之间也发生过迁移airflow.lineage.hook→airflow.sdk.lineage同样通过条件注册处理。回退算法create_module_getattr生成的__getattr__按以下顺序解析属性参见 _compat_utils.py 中__getattr__内部实现命中_RENAME_MAP先尝试新路径.新名字失败再尝试旧路径.旧名字两者都失败时抛出聚合了原始异常的ImportError报错信息形如Could not import X from new_path or old_name from old_path命中_MODULE_MAP按顺序importlib.import_module每个候选模块路径全部失败则抛ImportError命中_IMPORT_MAP对每个候选模块路径执行__import__(path, fromlist[name])后getattr任一成功即返回均未命中抛出标准AttributeError(fmodule has no attribute {name!r})。这种惰性 有序回退意味着导入airflow.providers.common.compat.sdk本身几乎没有开销只有真正访问某个属性时才会触发实际的模块导入与回退探测同时由于所有回退路径集中声明在模块级映射表中行为可预测、可静态检查。具体兼容场景剖析Dataset → Assetassets 子模块assets/init.py 展示了与sdk.py风格不同的硬分支写法运行期按AIRFLOW_V_3_0_PLUS直接二选一。Airflow 3 从airflow.sdk.definitions.asset导入Asset/AssetAlias/AssetAll/AssetAny从airflow.api_fastapi.auth.managers.models.resource_details导入AssetAliasDetails/AssetDetailsAirflow 2.x 则导入Dataset/DatasetAlias等并就地改名为 Asset 系列from airflow.datasets import Dataset as Asset同时把expand_alias_to_datasets别名为expand_alias_to_assets。对外暴露的统一名字在__all__中列出Asset、AssetAlias、AssetAliasDetails、AssetAll、AssetAny、AssetDetails、expand_alias_to_assets。TriggerBaseEventTrigger 的降级回退triggers.py 只处理一件事Airflow 3.0 新增了BaseEventTrigger用于标记可作为事件驱动调度源的触发器。Airflow 2.x 没有这个基类因此这里把它登记为_RENAME_MAP条目_RENAME_MAP { BaseEventTrigger: (airflow.triggers.base, airflow.triggers.base, BaseTrigger), }即 3.x 上解析为真正的BaseEventTrigger2.x 上回退为BaseTrigger——注释明确说明子类在 2.x 上仍可作为 deferred trigger 正常工作只是无法驱动 asset watcher。异步 Connection 与 Hook 的向后兼容助手connection/init.py 提供的get_async_connection(conn_id, hookNone)与 hook/init.py 提供的get_async_hook(conn_id, hook_paramsNone)是另一类兼容能力——为尚未原生支持异步的版本补齐async入口async def get_async_connection(conn_id: str, hookNone) - Connection: from asgiref.sync import sync_to_async hook hook or BaseHook if hasattr(hook, aget_connection): return await hook.aget_connection(conn_idconn_id) return await sync_to_async(hook.get_connection)(conn_idconn_id)其策略是先探测对象上是否存在aget_connection/aget_hook这类原生异步方法Airflow 3 的BaseHook具备存在则直接await否则用asgiref.sync.sync_to_async包装同步的get_connection/get_hook。这解释了为何asgiref是该包的必选依赖。hook参数支持传类或实例设计意图是让 Hook 子类对get_connection的重写也能被尊重。standard 子模块Operator 与 3.2 特性回退standard/operators.py 通过create_module_getattr统一导出BaseBranchOperator、BaseOperator、BranchMixIn、get_current_context等它们本身再走sdk的回退链以及PythonOperator、ShortCircuitOperator、_SERIALIZERSairflow.providers.standard.operators.python回退airflow.operators.python。其中对BaseAsyncOperator/is_async_callable的处理体现了版本不足时给出明确错误的工程取舍仅当AIRFLOW_V_3_2_PLUS为真时才从airflow.sdk导入真正的实现在 3.1 及以下版本上本文件现场定义一个桩类BaseAsyncOperator其execute会抛出RuntimeError(Async operators require Airflow 3.2. Upgrade Airflow or use a synchronous callable.)并本地实现了一个基于inspect.iscoroutinefunction含functools.partial解包的is_async_callable。其他兼容子模块包内还有若干面向特定领域的兼容模块均可在源码树中查证security/access_view.py、security/permissions.py跨版本的权限/访问视图相关兼容lineage/hook.py 与 lineage/entities.py血缘采集兼容sqlalchemy/orm.py跨版本 SQLAlchemy ORM 风格兼容openlineage/ 子目录check.py、facet.py、utils/spark.py、utils/sql.pyOpenLineage 兼容工具对应openlineageextra。测试如何验证兼容层单元测试 test_sdk.py 中的test_all_compat_imports_work遍历sdk.__all__中的每一个导出名断言至少一条回退路径可成功导入且非None并汇总所有失败项一次性pytest.fail——这保证了映射表新增条目后回退机制不会被静默破坏。test_invalid_import_raises_attribute_error则验证访问不存在的属性会抛出AttributeError消息含has no attribute NonExistentClass与__getattr__第 4 步的行为一一对应。同目录树下还有针对 connection、hook、lineage、openlineage、security、sqlalchemy 各兼容模块的对应测试文件如 test_hook.py可在 tests/unit/common/compat 下按模块逐一查看。小结common.compat以极小的依赖面仅apache-airflow2.11.0asgiref换来了跨大版本开发的核心能力写 Provider/Operator 时from airflow.providers.common.compat.sdk import BaseOperator, PythonOperator, ...一类导入可在 Airflow 2.x 与 3.x 上同时工作无需自己维护try/except ImportError需要条件分支时引入AIRFLOW_V_3_x_PLUS系列常量做显式版本判断需要异步能力时使用get_async_connection/get_async_hook获得带同步回退的async入口处理 Dataset/Asset 迁移时统一从common.compat.assets导入 Asset 家族符号。安装方式固定为pip install apache-airflow-providers-common-compat按需追加[openlineage]、[standard]extras最低要求 Airflow2.11.0、Python3.10与 pyproject.toml 的声明完全一致。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表