Repository navigation
[Data] Write Delta timestamps as microseconds in write_delta - #66907
HirokiNariyoshi wants to merge 5 commits into
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces automatic conversion of PyArrow timestamp types (including nested ones in structs, lists, and maps) to microsecond precision (and UTC for timezone-aware ones) when writing to Delta Lake, ensuring compatibility with Delta Lake's supported timestamp formats. The reviewer suggested extending this conversion to handle dictionary-encoded timestamp types to make the implementation more robust.
Signed-off-by: Hiroki Nariyoshi <hnariyos@uwaterloo.ca>
Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> Signed-off-by: Hiroki Nariyoshi <narihiro.hiroki@gmail.com> Signed-off-by: Hiroki Nariyoshi <hnariyos@uwaterloo.ca>
Signed-off-by: Hiroki Nariyoshi <hnariyos@uwaterloo.ca>
c404171 to
7b07e05
Compare
There was a problem hiding this comment.
Code Review
This pull request introduces automatic conversion of PyArrow timestamp types, including nested ones, to microsecond precision (and UTC for timezone-aware ones) when writing to Delta Lake. This ensures compatibility with Delta Lake's supported timestamp formats. Feedback points out a bug in the fixed-size list conversion logic where PyArrow would raise a ValueError if a Field object is passed with a specified list size, and suggests using the resolved DataType directly.
|
I've tested _cast_to_delta_schema() with an out-of-range timestamp[s] value (9223372036855), and it correctly raises ArrowInvalid instead of silently overflowing. Just sharing the result in case it is helpful. Since the docstring explicitly guarantees this behavior, would it be worth adding a regression test for this boundary case? It could help prevent accidental changes to the overflow behavior in future refactoring. Here's my attempt on that regression test Also, _to_delta_type() recursively converts both map keys and values, but the nested timestamp test only covers timestamp values with string keys. Would a timestamp-keyed map also be a supported case for this conversion? If so, could a test for that path be added, or clarify whether Delta Lake supports it? None of them are blockers, just suggestion for test coverage. |
…elta Signed-off-by: Hiroki Nariyoshi <hnariyos@uwaterloo.ca>
Thanks for testing this! Both addressed in d7a5a2e test_timestamp_conversion_rejects_out_of_range_values, based on your test, using 9,223,372,036,855 s, the smallest whole second whose microsecond value overflows int64. Map keys are converted, and timezone-aware keys commit fine, so I added a test for that. Naive timestamp keys can’t be written yet, independent of this PR. when a naive timestamp appears only as a map key, deltalake doesn’t add the timestampNtz feature, so the commit fails even with microsecond keys. write_deltalake fails the same way on 1.5.0 and 1.6.1. |
Description
DeltaDatasinkwrites Parquet with PyArrow directly and never convertstimestamps to microseconds, which deltalake's own
write_deltalakedoes.Any
s,ms, ornstimestamp, whether top-level, nested in a struct,list, or map, or dictionary-encoded, therefore breaks the write:
Invalid data type for Delta Lakeafter the workers have already writtentheir Parquet files, leaving those files orphaned.
nscommits a log schema ofusovernsfiles, so thetable's schema disagrees with its Parquet files. Timezone-aware
nscommits a table with the
timestampNanosreader feature, which thedeltalake reader rejects.
This PR casts every timestamp, including nested and dictionary-encoded ones,
to
timestamp[us](timezone-aware ones in UTC) on each worker beforewriting. It also converts the schema captured in
on_write_start, which iswhat an empty write commits. Sub-microsecond precision is truncated, since
Delta can't store it; out-of-range values still raise. Nanosecond timestamps are always
truncated to us, even if deltalake.enable_nanosecond_timestamps() (experimental, deltalake ≥ 1.6)
has been called, since Delta's stable timestamp type is microseconds.
The change is limited to
DeltaDatasinkand its tests.Repro on master with deltalake 1.5.0:
Related issues
None. I found this while reading
DeltaDatasink. I searched open and closedissues and PRs for
write_delta,DeltaDatasink,delta_datasink,timestampNanos, and "Invalid data type for Delta Lake", and found noduplicate.
Additional information
Tests run locally:
pytest python/ray/data/tests/datasource/test_write_delta.py: 112 passedon deltalake 1.5.0 and on 1.6.1.
large_list, fixed_size_list, map, and dictionary-encoded columns; an
nsappend onto a
ustable; an empty dataset; out-of-range values; timestamp map keys)fail on master and pass with this change.
file allows.
pre-commit run --files <changed files>: pass.AI assistance was used to draft the change.
I reviewed every changed line and ran the tests above.