Apache Airflow SDK Temporal Partition Mapper 接口变更:timezone 参数支持与 Keyword-Only 构造器迁移指南
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本篇基于 Airflow 仓库的变更说明 67164.significant.rst,讲解airflow.sdk中六个 Temporal Partition Mapper 类(StartOfHourMapper、StartOfDayMapper、StartOfWeekMapper、StartOfMonthMapper、StartOfQuarterMapper、StartOfYearMapper)的构造器接口变更:新增与 core 对齐的timezone关键字参数、构造器全面改为 keyword-only,以及字符串时区在构造时即通过parse_timezone解析并提前报错的行为变化。读完后你将能安全地把存量StartOf*Mapper(...)调用点迁移到新接口,并理解其背后的源码实现与序列化语义。
一、变更背景:SDK 与 core 的 Temporal Mapper 对齐
在 Airflow 的资产(Asset)分区机制中,Partition Mapper 负责把上游资产事件的分区键映射为触发下游 Dag 运行的目标分区键。Temporal 系列 Mapper 是其中最常用的一类:它们把一个时间戳键归一化到其所属周期(小时/天/周/月/季/年)的起点,例如2024-03-13T10:42:15→2024-03-13(按天)。
core 侧的实现位于 airflow-core/src/airflow/partition_mappers/temporal.py,SDK 侧(即用户写 Dag 时从airflow.sdk导入的入口)位于 task-sdk/src/airflow/sdk/definitions/partition_mappers/temporal.py。这两套类层级相互独立——SDK 不能导入 core,因此两边各自维护同一份默认值与语义。本次变更(newsfragment 67164)的核心目的就是消除二者签名上的历史偏差:
StartOfHourMapper、StartOfDayMapper、StartOfWeekMapper、StartOfMonthMapper、StartOfQuarterMapper、StartOfYearMapper(从airflow.sdk导入)现在都接受timezone关键字参数,与 core 的_BaseTemporalMapper签名一致;- 构造器改为 keyword-only,
input_format和output_format不再接受位置参数。
二、行为变化详解
2.1input_format/output_format不再接受位置参数
在 task-sdk 1.2.1 中,下面的写法是合法的(位置传参):
StartOfDayMapper("%Y-%m-%dT%H:%M:%S")新接口下,必须改为按名称传参:
StartOfDayMapper(input_format="%Y-%m-%dT%H:%M:%S")迁移方法很直接:检查所有StartOf*Mapper(...)调用点,把位置实参改为关键字实参。变更说明中给出的 Migration 建议也是这一条。
2.2 字符串timezone在构造时立即解析
此前如果传入一个无法识别的时区名称字符串,它会被原样存储,问题延迟到运行后期才暴露,甚至在某些序列化路径中被静默丢弃。现在字符串timezone在构造阶段就通过parse_timezone解析,未知名称会立即抛出pendulum.tz.exceptions.InvalidTimezone。
parse_timezone的实现见 task-sdk/src/airflow/sdk/_shared/timezones/timezone.py:其接口与pendulum.timezone(name)相同,接受 IANA 时区名或相对 UTC 的偏移秒数;未知名称即由 pendulum 抛出InvalidTimezone。这意味着时区拼写错误(例如把America/New_York写成America/Newyork)会在 Dag 解析期就 fail-fast,而不是留到调度器反序列化或序列化落库时才发现。
三、新构造器签名与参数说明
从 SDK 源码 task-sdk/src/airflow/sdk/definitions/partition_mappers/temporal.py 看,_BaseTemporalMapper基于attrs定义,三个字段全部声明为kw_only=True:
@attrs.define class _BaseTemporalMapper(PartitionMapper): """Base class for Temporal Partition Mappers.""" default_output_format: ClassVar[str] expected_decoded_type: ClassVar[type] = datetime _timezone: str | Timezone | FixedTimezone = attrs.field( alias="timezone", default="UTC", kw_only=True, converter=_timezone_converter, ) input_format: str = attrs.field(default="%Y-%m-%dT%H:%M:%S", kw_only=True) output_format: str | None = attrs.field(default=None, kw_only=True) def __attrs_post_init__(self) -> None: if not self.output_format: self.output_format = self.default_output_format参数语义如下(与 core 侧 airflow-core/src/airflow/partition_mappers/temporal.py 的__init__签名一一对应):
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
timezone | str/pendulum.Timezone/pendulum.FixedTimezone | "UTC" | 用于本地化 naive 上游键的时区。字符串在构造时经_timezone_converter→parse_timezone解析;未知名称立即抛InvalidTimezone |
input_format | str | "%Y-%m-%dT%H:%M:%S" | 解析上游分区键所用的strptime兼容格式 |
output_format | str \| None | None(实际取各子类的default_output_format) | 生成下游(周期起点)键所用的strftime格式;为None或空时回退到子类的default_output_format |
另外,core 侧基类构造器还带一个max_downstream_keys: int | None = None参数(见 airflow-core/src/airflow/partition_mappers/temporal.py),用于约束 fan-out 等场景生成的下游键数量上限,同样只能按关键字传入。
六个子类各自的默认output_format(SDK temporal.py):
| Mapper 类 | 默认output_format | 映射示例 |
|---|---|---|
StartOfHourMapper | %Y-%m-%dT%H | 2024-03-13T10:42:15→2024-03-13T10 |
StartOfDayMapper | %Y-%m-%d | 2024-03-13T10:42:15→2024-03-13 |
StartOfWeekMapper | %Y-%m-%d (W%V) | 2024-03-13T10:42:15→2024-03-11 (W11)(ISO 周周一) |
StartOfMonthMapper | %Y-%m | 2024-03-13T10:42:15→2024-03 |
StartOfQuarterMapper | %Y-Q{quarter} | 2024-03-13T10:42:15→2024-Q1 |
StartOfYearMapper | %Y | 2024-03-13T10:42:15→2024 |
四、源码级印证:timezone如何参与映射与序列化
4.1 时区参与键的解析与归一化
core 侧_BaseTemporalMapper.to_downstream(temporal.py)展示了timezone的运行时作用:
def to_downstream(self, key: str) -> str: dt = datetime.strptime(key, self.input_format) if dt.tzinfo is None: dt = make_aware(dt, self._timezone) else: dt = dt.astimezone(self._timezone) normalized = self.normalize(dt) return self.format(normalized)即:naive 上游键先按 mapper 的timezone本地化,aware 键则转换到该时区,然后才归一化到周期起点并格式化。to_partition_date、encode_upstream同样使用该时区,保证RollupMapper等组合场景下“解码-编码”往返后时刻一致。这正是 SDK 过去缺失timezone参数会造成语义偏差的地方——现在两侧签名一致后,SDK 侧的 Dag 定义与 core 侧反序列化后的行为可以精确对齐。
4.2 序列化路径:从“静默丢弃”到“显式编码”
core 侧serialize(temporal.py)会把时区显式编码进序列化字典:
def serialize(self) -> dict[str, Any]: from airflow.serialization.encoders import encode_timezone result: dict[str, Any] = { "timezone": encode_timezone(self._timezone), "input_format": self.input_format, "output_format": self.output_format, } ...反序列化时deserialize通过cls(timezone=parse_timezone(...), ...)按关键字重建。newsfragment 中提到的“某些路径下被静默丢弃”的问题,正是因为旧实现没有构造期校验与统一编码契约:现在字符串时区在__attrs_post_init__之前的 converter 里就被解析成Timezone/FixedTimezone对象,序列化与反序列化两端都以解析后的对象为准,消除了不一致窗口。
4.3 构造期 fail-fast 的连带好处
core 侧StartOfWeekMapper与StartOfQuarterMapper还在构造时预先编译output_format对应的正则(_compile_output_format_regex),格式非法会在 Dag 解析期抛ValueError,而不是在调度器 tick 深处抛出难以定位的re.error(temporal.py)。这与“时区未知即报错”的设计取向一致:所有可验证的配置错误都尽量前移到构造点。
五、迁移步骤与用法示例
- 全局检索:在所有 Dag 目录中检索
StartOfHourMapper(、StartOfDayMapper(、StartOfWeekMapper(、StartOfMonthMapper(、StartOfQuarterMapper(、StartOfYearMapper(的调用点; - 改位置传参为关键字传参:
StartOfDayMapper("%Y-%m-%dT%H:%M:%S")→StartOfDayMapper(input_format="%Y-%m-%dT%H:%M:%S"); - 按需显式指定时区:若上游资产事件的键是某业务时区的 wall-clock 时间,传入
timezone="America/New_York"等 IANA 名称或pendulum.Timezone/FixedTimezone对象;注意拼写错误现在会直接导致 Dag 解析失败(这是预期行为,属于 fail-fast); - 验证:本地解析 Dag,确认不再出现
TypeError(位置传参)或InvalidTimezone。
完整用法可以参考示例 DAG airflow-core/src/airflow/example_dags/example_asset_partition.py,其中涵盖了多资产分区对齐(StartOfHourMapper)、滚动聚合(RollupMapper(upstream_mapper=StartOfDayMapper(), ...))、月粒度聚合(StartOfMonthMapper(input_format="%Y-%m-%d"))以及按周扇出(FanOutMapper(upstream_mapper=StartOfWeekMapper(), window=WeekWindow()))等典型模式。一个最小的迁移前后对照:
# 迁移前(task-sdk 1.2.1 及更早,位置传参) mapper = StartOfDayMapper("%Y-%m-%dT%H:%M:%S") # 迁移后(keyword-only,可选指定时区) mapper = StartOfDayMapper( input_format="%Y-%m-%dT%H:%M:%S", timezone="Asia/Shanghai", )行为验证可以参考 SDK 侧测试 task-sdk/tests/task_sdk/definitions/test_partition_mappers.py,其中TestSdkCategoricalRollupGuard等用例覆盖了StartOfDayMapper与DayWindow的类型配对校验;core 侧对应测试为 airflow-core/tests/unit/partition_mappers/test_temporal.py。
六、小结与注意事项
- 本次变更属于Behaviour changes与Code interface changes两类(见 newsfragment 末尾的变更类型勾选),对依赖 task-sdk 1.2.1 位置传参写法的存量代码是不兼容变更;
- 六个 Temporal Mapper 的构造器现在与 core 的
_BaseTemporalMapper签名对齐:timezone(默认"UTC")、input_format、output_format均为 keyword-only; - 字符串时区在构造期解析,未知时区名立即抛
pendulum.tz.exceptions.InvalidTimezone,不再可能“存着错值、后期才炸或静默丢失”; - 适用前提:以上结论基于当前仓库中 task-sdk 与 airflow-core 的源码;如果你固定在 task-sdk 1.2.1 及其更早版本,位置传参写法仍可用,建议尽快按上文迁移,以便获得与 core 完全一致的时区语义与序列化行为。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考