Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions docs/docs/pypaimon/python-api.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -1119,8 +1119,9 @@ sub-columns for column-skipping via sub-field projection.
Fields not listed in `variant.shreddingSchema` are stored in the overflow `value` bytes and remain
fully accessible on the read path.

Supported Paimon type strings for shredded sub-fields: `BOOLEAN`, `INT`, `BIGINT`, `FLOAT`, `DOUBLE`,
`VARCHAR`, `DECIMAL(p,s)`, and nested `ROW` types for recursive object shredding.
Supported Paimon type strings for shredded sub-fields: `BOOLEAN`, `TINYINT`, `SMALLINT`, `INT`,
`BIGINT`, `FLOAT`, `DOUBLE`, `VARCHAR`, `BINARY`, `VARBINARY`, `DECIMAL(p,s)`, `ARRAY`, and nested
`ROW` types for recursive object shredding.

</TabItem>

Expand Down
1 change: 1 addition & 0 deletions docs/sidebars.js
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ const sidebars = {
},
"items": [
"multimodal-table/data-evolution",
"multimodal-table/variant",
"multimodal-table/blob",
"multimodal-table/vector",
{
Expand Down
64 changes: 61 additions & 3 deletions paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.data.variant.GenericVariant;
import org.apache.paimon.data.variant.GenericVariantBuilder;
import org.apache.paimon.data.variant.GenericVariantUtil.Type;
import org.apache.paimon.deletionvectors.BitmapDeletionVector;
import org.apache.paimon.deletionvectors.DeletionVector;
import org.apache.paimon.deletionvectors.append.BaseAppendDeleteFileMaintainer;
Expand Down Expand Up @@ -101,6 +103,10 @@
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
Expand Down Expand Up @@ -1682,6 +1688,28 @@ public void testJavaWriteVariantTable() throws Exception {
3,
BinaryString.fromString("Carol"),
GenericVariant.fromJson("[1,2,3]")));

// Scalar DATE/TIMESTAMP/TIMESTAMP_NTZ/UUID values for cross-language compatibility.
GenericVariantBuilder b = new GenericVariantBuilder(false);
b.appendDate((int) LocalDate.of(2024, 1, 15).toEpochDay());
GenericVariant v = b.result();
write.write(GenericRow.of(4, BinaryString.fromString("Dave"), v));

b = new GenericVariantBuilder(false);
b.appendTimestampNtz(
toEpochMicros(LocalDateTime.of(2024, 1, 15, 12, 30, 45, 123456000)));
v = b.result();
write.write(GenericRow.of(5, BinaryString.fromString("Eve"), v));

b = new GenericVariantBuilder(false);
b.appendTimestamp(toEpochMicros(Instant.parse("2024-01-15T12:30:45.123456Z")));
v = b.result();
write.write(GenericRow.of(6, BinaryString.fromString("Frank"), v));

b = new GenericVariantBuilder(false);
b.appendUuid(UUID.fromString("12345678-1234-5678-1234-567812345678"));
v = b.result();
write.write(GenericRow.of(7, BinaryString.fromString("Grace"), v));
commit.commit(write.prepareCommit());
}

Expand All @@ -1691,7 +1719,7 @@ public void testJavaWriteVariantTable() throws Exception {
TableRead read = readTable.newRead();
List<String> res =
getResult(read, splits, row -> internalRowToString(row, readTable.rowType()));
assertThat(res).hasSize(3);
assertThat(res).hasSize(7);
LOG.info("testJavaWriteVariantTable: wrote and read back {} VARIANT rows", res.size());

// Also write a shredded VARIANT table for Python to read (variant_shredded_test).
Expand Down Expand Up @@ -1762,7 +1790,7 @@ public void testJavaReadVariantTable() throws Exception {
TableRead read = table.newRead();
List<String> res =
getResult(read, splits, row -> internalRowToString(row, table.rowType()));
assertThat(res).hasSize(4);
assertThat(res).hasSize(7);

// Verify the VARIANT column is present in the schema
assertThat(table.rowType().getFieldNames()).contains("payload");
Expand All @@ -1781,8 +1809,29 @@ public void testJavaReadVariantTable() throws Exception {
assertThat(row.isNullAt(2)).isTrue();
} else {
assertThat(row.isNullAt(2)).isFalse();
org.apache.paimon.data.variant.Variant v = row.getVariant(2);
GenericVariant v = (GenericVariant) row.getVariant(2);
assertThat(v).isNotNull();
if (id == 5) {
// DATE '2024-01-15'
assertThat(v.getType()).isEqualTo(Type.DATE);
assertThat(v.getLong())
.isEqualTo(LocalDate.of(2024, 1, 15).toEpochDay());
} else if (id == 6) {
// TIMESTAMP_NTZ '2024-01-15 12:30:45.123456'
assertThat(v.getType()).isEqualTo(Type.TIMESTAMP_NTZ);
long expectedMicros =
toEpochMicros(
LocalDateTime.of(
2024, 1, 15, 12, 30, 45, 123456000));
assertThat(v.getLong()).isEqualTo(expectedMicros);
} else if (id == 7) {
// UUID '12345678-1234-5678-1234-567812345678'
assertThat(v.getType()).isEqualTo(Type.UUID);
assertThat(v.getUuid())
.isEqualTo(
UUID.fromString(
"12345678-1234-5678-1234-567812345678"));
}
}
});
}
Expand Down Expand Up @@ -1830,6 +1879,15 @@ public void testJavaReadVariantTable() throws Exception {
shreddedRes.size());
}

private static long toEpochMicros(LocalDateTime dateTime) {
return dateTime.toInstant(ZoneOffset.UTC).getEpochSecond() * 1_000_000L
+ dateTime.getNano() / 1000L;
}

private static long toEpochMicros(Instant instant) {
return instant.getEpochSecond() * 1_000_000L + instant.getNano() / 1000L;
}

/** Step 1: Write 5 base files for compact conflict test. */
@Test
@EnabledIfSystemProperty(named = "run.e2e.tests", matches = "true")
Expand Down
33 changes: 30 additions & 3 deletions paimon-python/pypaimon/data/generic_variant.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
v.metadata() – raw metadata bytes
"""

import calendar
import datetime
import decimal as _decimal
import enum
Expand Down Expand Up @@ -353,6 +354,13 @@ def append_binary(self, b):
self._buf[self._pos:self._pos + len(b)] = b
self._pos += len(b)

def append_uuid(self, u):
# UUID values are 16-byte big-endian: msb followed by lsb.
self._write_byte(_primitive_header(_UUID))
self._ensure(16)
self._buf[self._pos:self._pos + 16] = u.bytes
self._pos += 16

def append_date(self, days_since_epoch):
self._write_byte(_primitive_header(_DATE))
self._write_le(days_since_epoch & 0xFFFFFFFF, 4)
Expand Down Expand Up @@ -447,9 +455,28 @@ def build_python(self, obj):
self._finish_writing_array(start, elem_offsets)
elif isinstance(obj, bytes):
self.append_binary(obj)
elif isinstance(obj, _uuid.UUID):
self.append_uuid(obj)
elif isinstance(obj, datetime.datetime):
micros = self._datetime_to_micros(obj)
if obj.tzinfo is not None:
self.append_timestamp(micros)
else:
self.append_timestamp_ntz(micros)
elif isinstance(obj, datetime.date):
days = (obj - _EPOCH_DATE).days
self.append_date(days)
else:
raise TypeError(f'Unsupported Python type for variant encoding: {type(obj).__name__}')

@staticmethod
def _datetime_to_micros(dt):
"""Convert a datetime to microseconds since epoch using pure integer arithmetic. """
if dt.tzinfo is not None:
dt = dt.astimezone(datetime.timezone.utc)
seconds = calendar.timegm(dt.timetuple())
return seconds * 1_000_000 + dt.microsecond

def _try_decimal_or_double(self, d):
try:
sign, digits, exponent = d.as_tuple()
Expand Down Expand Up @@ -687,9 +714,9 @@ def _to_python_impl(self, value, metadata, pos):
length = _read_unsigned(value, pos + 1, _U32_SIZE)
return bytes(value[pos + 1 + _U32_SIZE:pos + 1 + _U32_SIZE + length])
if vtype == _Type.UUID:
# 16 bytes: two little-endian int64 (msb, lsb) → standard UUID
msb = _read_unsigned(value, pos + 1, 8)
lsb = _read_unsigned(value, pos + 9, 8)
# UUID values are 16-byte big-endian: msb followed by lsb.
msb = int.from_bytes(value[pos + 1:pos + 9], 'big', signed=False)
lsb = int.from_bytes(value[pos + 9:pos + 17], 'big', signed=False)
return _uuid.UUID(int=(msb << 64) | lsb)
if vtype == _Type.OBJECT:
def _build_dict(size, id_size, offset_size, id_start, offset_start, data_start):
Expand Down
46 changes: 42 additions & 4 deletions paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import os
import sys
import unittest
import uuid
from decimal import Decimal

import pandas as pd
Expand Down Expand Up @@ -1931,7 +1932,7 @@ def test_py_read_variant_table(self):
splits = table_scan.plan().splits()
result = table_read.to_arrow(splits)

self.assertEqual(result.num_rows, 3)
self.assertEqual(result.num_rows, 7)

# VARIANT maps to struct<value: binary NOT NULL, metadata: binary NOT NULL>
payload_field = result.schema.field('payload')
Expand Down Expand Up @@ -1973,6 +1974,31 @@ def test_py_read_variant_table(self):
carol_data = GenericVariant.from_arrow_struct(payload_list[id_list.index(3)]).to_python()
self.assertEqual(carol_data, [1, 2, 3])

# Row 4: Dave, DATE '2024-01-15'
dave_data = GenericVariant.from_arrow_struct(payload_list[id_list.index(4)]).to_python()
self.assertEqual(dave_data, datetime.date(2024, 1, 15))

# Row 5: Eve, TIMESTAMP_NTZ '2024-01-15 12:30:45.123456'
eve_data = GenericVariant.from_arrow_struct(payload_list[id_list.index(5)]).to_python()
self.assertEqual(
eve_data, datetime.datetime(2024, 1, 15, 12, 30, 45, 123456)
)

# Row 6: Frank, TIMESTAMP '2024-01-15 12:30:45.123456 UTC'
frank_data = GenericVariant.from_arrow_struct(payload_list[id_list.index(6)]).to_python()
self.assertEqual(
frank_data,
datetime.datetime(
2024, 1, 15, 12, 30, 45, 123456, tzinfo=datetime.timezone.utc
),
)

# Row 7: Grace, UUID '12345678-1234-5678-1234-567812345678'
grace_data = GenericVariant.from_arrow_struct(payload_list[id_list.index(7)]).to_python()
self.assertEqual(
grace_data, uuid.UUID('12345678-1234-5678-1234-567812345678')
)

print("test_py_read_variant_table: verified {} VARIANT rows".format(result.num_rows))

# Also verify shredded VARIANT: Java wrote variant_shredded_test with
Expand Down Expand Up @@ -2029,6 +2055,9 @@ def test_py_write_variant_table(self):
id=2 payload=[10,20,30]
id=3 payload="hello"
id=4 payload=null
id=5 payload=DATE '2024-01-15'
id=6 payload=TIMESTAMP_NTZ '2024-01-15 12:30:45.123456'
id=7 payload=UUID '12345678-1234-5678-1234-567812345678'
"""
variant_type = pa.struct([
pa.field('value', pa.binary(), nullable=False),
Expand All @@ -2046,15 +2075,24 @@ def test_py_write_variant_table(self):
self.catalog.create_table(table_name, schema, False)
table = self.catalog.get_table(table_name)

test_uuid = uuid.UUID('12345678-1234-5678-1234-567812345678')
variant_col = GenericVariant.to_arrow_array([
GenericVariant.from_python({"name": "test", "value": 42}),
GenericVariant.from_python([10, 20, 30]),
GenericVariant.from_python("hello"),
None, # SQL NULL at the column level, not a VARIANT containing JSON null
GenericVariant.from_python(datetime.date(2024, 1, 15)),
GenericVariant.from_python(
datetime.datetime(2024, 1, 15, 12, 30, 45, 123456)
),
GenericVariant.from_python(test_uuid),
])
data = pa.table({
'id': pa.array([1, 2, 3, 4], type=pa.int32()),
'name': pa.array(['row1', 'row2', 'row3', 'row4'], type=pa.string()),
'id': pa.array([1, 2, 3, 4, 5, 6, 7], type=pa.int32()),
'name': pa.array(
['row1', 'row2', 'row3', 'row4', 'row5', 'row6', 'row7'],
type=pa.string()
),
'payload': variant_col,
}, schema=pa_schema)

Expand All @@ -2065,7 +2103,7 @@ def test_py_write_variant_table(self):
table_commit.commit(table_write.prepare_commit())
table_write.close()
table_commit.close()
print("test_py_write_variant_table: wrote 4 VARIANT rows to {}".format(table_name))
print("test_py_write_variant_table: wrote 7 VARIANT rows to {}".format(table_name))

# Also write a shredded VARIANT table (py_variant_shredded_test) for Java to read.
# Python shreds the 'age' (BIGINT) and 'city' (VARCHAR) sub-fields of 'payload'
Expand Down
Loading
Loading