ARTICLE DETAIL

资讯详情

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

Apache Airflow SDK Temporal Partition Mapper 接口变更:timezone 参数支持与 Keyword-Only 构造器迁移指南

Apache Airflow SDK Temporal Partition Mapper 接口变更:timezone 参数支持与 Keyword-Only 构造器迁移指南 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.pySDK 侧即用户写 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-onlyinput_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_onlyTrueattrs.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( aliastimezone, defaultUTC, kw_onlyTrue, converter_timezone_converter, ) input_format: str attrs.field(default%Y-%m-%dT%H:%M:%S, kw_onlyTrue) output_format: str | None attrs.field(defaultNone, kw_onlyTrue) 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__签名一一对应参数类型默认值说明timezonestr/pendulum.Timezone/pendulum.FixedTimezoneUTC用于本地化 naive 上游键的时区。字符串在构造时经_timezone_converter→parse_timezone解析未知名称立即抛InvalidTimezoneinput_formatstr%Y-%m-%dT%H:%M:%S解析上游分区键所用的strptime兼容格式output_formatstr \| NoneNone实际取各子类的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_formatSDK temporal.pyMapper 类默认output_format映射示例StartOfHourMapper%Y-%m-%dT%H2024-03-13T10:42:15→2024-03-13T10StartOfDayMapper%Y-%m-%d2024-03-13T10:42:15→2024-03-13StartOfWeekMapper%Y-%m-%d (W%V)2024-03-13T10:42:15→2024-03-11 (W11)ISO 周周一StartOfMonthMapper%Y-%m2024-03-13T10:42:15→2024-03StartOfQuarterMapper%Y-Q{quarter}2024-03-13T10:42:15→2024-Q1StartOfYearMapper%Y2024-03-13T10:42:15→2024四、源码级印证timezone如何参与映射与序列化4.1 时区参与键的解析与归一化core 侧_BaseTemporalMapper.to_downstreamtemporal.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 侧serializetemporal.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(timezoneparse_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.errortemporal.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 时间传入timezoneAmerica/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_mapperStartOfDayMapper(), ...)、月粒度聚合StartOfMonthMapper(input_format%Y-%m-%d)以及按周扇出FanOutMapper(upstream_mapperStartOfWeekMapper(), windowWeekWindow())等典型模式。一个最小的迁移前后对照# 迁移前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, timezoneAsia/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),仅供参考
返回列表