From dd61e06a345f2a7453c50b3991e9d57078845ee2 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Wed, 23 Sep 2026 02:03:53 +0800 Subject: [PATCH] [python] Read row-format TIMESTAMP values in the precision's time unit The row-format reader returned a TIMESTAMP value in milliseconds for precision <= 3 and microseconds for precision > 3, but PyarrowFieldParser.from_paimon_type maps the precision to four Arrow units (0 -> s, 1-3 -> ms, 4-6 -> us, 7-9 -> ns) and the value is placed into that type without conversion. Only precisions 1-6 agreed: a precision-0 value was read as seconds from a millisecond integer (x1000, overflowing to year 52626), and a precision 7-9 value was read as nanoseconds from a microsecond integer (/1000). Return the value in the Arrow unit the precision maps to: seconds for 0, millis for 1-3, micros for 4-6, nanos for 7-9. The wire format (a millis long plus, only for precision > 3, a nano_of_milli varint) is unchanged and already written that way, so this is a read-side fix. --- .../pypaimon/read/reader/format_row_reader.py | 13 ++-- .../tests/test_format_row_reader_writer.py | 61 +++++++++++++++++++ 2 files changed, 70 insertions(+), 4 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/format_row_reader.py b/paimon-python/pypaimon/read/reader/format_row_reader.py index 21b7e75c3fe0..d65355310db2 100644 --- a/paimon-python/pypaimon/read/reader/format_row_reader.py +++ b/paimon-python/pypaimon/read/reader/format_row_reader.py @@ -495,12 +495,17 @@ def _read_field(decoder: _RowDecoder, data_type) -> Any: elif type_name.startswith('TIMESTAMP'): precision = _parse_timestamp_precision(type_name) millis = decoder.read_long() + # The value is placed into the Arrow time unit from_paimon_type maps the + # precision to (0 -> s, 1-3 -> ms, 4-6 -> us, 7-9 -> ns), so it must be + # returned in that unit. nano_of_milli is only on the wire for precision > 3. + if precision == 0: + return millis // 1000 if precision <= 3: return millis - else: - nano_of_milli = decoder.read_var_int() - micros = millis * 1000 + nano_of_milli // 1000 - return micros + nano_of_milli = decoder.read_var_int() + if precision <= 6: + return millis * 1000 + nano_of_milli // 1000 + return millis * 1_000_000 + nano_of_milli elif type_name == 'VARIANT': value_bytes = decoder.read_bytes() metadata_bytes = decoder.read_bytes() diff --git a/paimon-python/pypaimon/tests/test_format_row_reader_writer.py b/paimon-python/pypaimon/tests/test_format_row_reader_writer.py index d8b5d9b7d190..f082947b148d 100644 --- a/paimon-python/pypaimon/tests/test_format_row_reader_writer.py +++ b/paimon-python/pypaimon/tests/test_format_row_reader_writer.py @@ -15,6 +15,7 @@ # specific language governing permissions and limitations # under the License. +import datetime import os import struct import tempfile @@ -150,6 +151,66 @@ def _varint(x): got = _read_field(_RowDecoder(buf, 0), AtomicType("DECIMAL(38, 10)")) assert got == Decimal("1234567890123456789012345678.9012345678") + def test_timestamp_precisions(self): + # from_paimon_type maps the precision to an Arrow unit (0 -> s, 1-3 -> ms, + # 4-6 -> us, 7-9 -> ns); the reader must return the value in that unit. + # Precision 0 overflowed (millis read as seconds) and 7-9 read low (micros + # read as nanos); 1-6 were already correct. + fields = [ + DataField(0, "ts0", AtomicType("TIMESTAMP(0)")), + DataField(1, "ts3", AtomicType("TIMESTAMP(3)")), + DataField(2, "ts6", AtomicType("TIMESTAMP(6)")), + DataField(3, "ts9", AtomicType("TIMESTAMP(9)")), + ] + base = datetime.datetime(2020, 9, 13, 12, 26, 40) + ts0 = base + ts3 = base.replace(microsecond=123000) + ts6 = base.replace(microsecond=123456) + ts9 = base.replace(microsecond=123456) + data = pa.table({ + "ts0": pa.array([ts0], type=pa.timestamp('s')), + "ts3": pa.array([ts3], type=pa.timestamp('ms')), + "ts6": pa.array([ts6], type=pa.timestamp('us')), + "ts9": pa.array([ts9], type=pa.timestamp('ns')), + }) + + with tempfile.NamedTemporaryFile(suffix=".row", delete=False) as tmp: + path = tmp.name + + try: + _write_row_file(path, fields, data) + result = _read_row_file(path, fields) + assert result.column("ts0").to_pylist() == [ts0] + assert result.column("ts3").to_pylist() == [ts3] + assert result.column("ts6").to_pylist() == [ts6] + assert result.column("ts9").to_pylist() == [ts9] + finally: + os.unlink(path) + + def test_timestamp_nanos_decoded_from_wire(self): + # A Java-written TIMESTAMP(9) carries nano_of_milli in 0..999999 (genuine + # sub-millisecond nanoseconds). The Python writer only emits multiples of + # 1000, so decode a hand-built wire buffer to exercise the ns formula with a + # non-multiple-of-1000 nano_of_milli directly. + from pypaimon.read.reader.format_row_reader import _read_field, _RowDecoder + + def _varint(n): + out = bytearray() + while True: + b = n & 0x7F + n >>= 7 + if n: + out.append(b | 0x80) + else: + out.append(b) + return bytes(out) + + millis = 1600000000000 + nano_of_milli = 123456 + buf = struct.pack('