ColtenOuO commented on code in PR #73748:
URL: https://github.com/apache/airflow/pull/73748#discussion_r4114361288
##########
task-sdk/src/airflow/sdk/serde/serializers/datetime.py:
##########
@@ -93,6 +103,12 @@ def deserialize(cls: type, version: int, data: dict | str)
-> datetime.date | da
)
if cls is datetime.datetime and isinstance(data, dict):
+ if tz is None:
+ # No timezone was stored, so this was a naive datetime. Interpret
the
+ # epoch in the configured default timezone (mirroring
``serialize``)
+ # and return a naive datetime, so the round-trip compares equal to
+ # what was pushed regardless of the deserializing process's OS
timezone.
Review Comment:
Nit: these two comment blocks restate each other and mostly narrate the
code. (L56-60)
##########
task-sdk/src/airflow/sdk/serde/serializers/datetime.py:
##########
@@ -93,6 +103,12 @@ def deserialize(cls: type, version: int, data: dict | str)
-> datetime.date | da
)
if cls is datetime.datetime and isinstance(data, dict):
+ if tz is None:
+ # No timezone was stored, so this was a naive datetime. Interpret
the
+ # epoch in the configured default timezone (mirroring
``serialize``)
+ # and return a naive datetime, so the round-trip compares equal to
+ # what was pushed regardless of the deserializing process's OS
timezone.
+ return
make_naive(datetime.datetime.fromtimestamp(float(data[TIMESTAMP]),
tz=datetime.timezone.utc))
Review Comment:
This line is 112 characters, over the project's `line-length = 110`, which
is likely why Static checks is red. `prek run ruff-format --from-ref main`
should fix.
##########
task-sdk/src/airflow/sdk/serde/serializers/datetime.py:
##########
@@ -52,9 +52,19 @@ def serialize(o: object) -> tuple[U, str, int, bool]:
if isinstance(o, datetime):
qn = qualname(o)
- tz = serialize_timezone(o.tzinfo) if o.tzinfo else None
+ if o.tzinfo is None:
+ # A naive datetime carries no timezone information, so anchor it to
+ # the configured default timezone (``core.default_timezone``)
rather
+ # than the OS local timezone of the serializing process. Otherwise
+ # the stored epoch silently depends on which machine writes the
value.
+ # The payload keeps ``tz`` empty so it still deserializes as naive.
+ ts = make_aware(o).timestamp()
+ tz = None
+ else:
+ ts = o.timestamp()
+ tz = serialize_timezone(o.tzinfo)
- return {TIMESTAMP: o.timestamp(), TIMEZONE: tz}, qn, __version__, True
+ return {TIMESTAMP: ts, TIMEZONE: tz}, qn, __version__, True
Review Comment:
The meaning of `timestamp` for naive values changes here (writer's OS zone →
`default_timezone`), but `__version__` (L34) stays at 2, so the reader can't
tell old payloads from new ones. This isn't limited to XCom: trigger kwargs
(`airflow/models/trigger.py`) and the XCom JSON encoder/decoder in core go
through the same serde, and old and new workers will run side by side during a
rolling upgrade.
Would it make sense to bump to `__version__ = 3` and apply the new
interpretation only for `version >= 3`, keeping the previous behaviour for
v1/v2 naive payloads? That way reads of existing data stay predictable instead
of shifting silently.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]