Skip to content

Commit 8b8585f

Browse files
committed
GH-3609: Add new sort order for int96 timestamps
1 parent 07812b9 commit 8b8585f

8 files changed

Lines changed: 490 additions & 26 deletions

File tree

parquet-column/src/main/java/org/apache/parquet/schema/ColumnOrder.java

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,11 +41,17 @@ public enum ColumnOrderName {
4141
* The column order is defined by the IEEE 754 standard.
4242
*/
4343
IEEE_754_TOTAL_ORDER,
44+
/**
45+
* Chronological order for INT96 timestamps.
46+
*/
47+
INT96_TIMESTAMP_ORDER
4448
}
4549

4650
private static final ColumnOrder UNDEFINED_COLUMN_ORDER = new ColumnOrder(ColumnOrderName.UNDEFINED);
4751
private static final ColumnOrder TYPE_DEFINED_COLUMN_ORDER = new ColumnOrder(ColumnOrderName.TYPE_DEFINED_ORDER);
4852
private static final ColumnOrder IEEE_754_TOTAL_ORDER = new ColumnOrder(ColumnOrderName.IEEE_754_TOTAL_ORDER);
53+
private static final ColumnOrder INT96_TIMESTAMP_COLUMN_ORDER =
54+
new ColumnOrder(ColumnOrderName.INT96_TIMESTAMP_ORDER);
4955

5056
/**
5157
* @return a {@link ColumnOrder} instance representing an undefined order
@@ -71,6 +77,14 @@ public static ColumnOrder ieee754TotalOrder() {
7177
return IEEE_754_TOTAL_ORDER;
7278
}
7379

80+
/**
81+
* @return a {@link ColumnOrder} instance representing the chronological order of INT96 timestamps
82+
* @see ColumnOrderName#INT96_TIMESTAMP_ORDER
83+
*/
84+
public static ColumnOrder int96TimestampOrder() {
85+
return INT96_TIMESTAMP_COLUMN_ORDER;
86+
}
87+
7488
private final ColumnOrderName columnOrderName;
7589

7690
private ColumnOrder(ColumnOrderName columnOrderName) {

parquet-column/src/main/java/org/apache/parquet/schema/PrimitiveComparator.java

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
import java.io.Serializable;
2222
import java.nio.ByteBuffer;
23+
import java.nio.ByteOrder;
2324
import java.util.Comparator;
2425
import org.apache.parquet.io.api.Binary;
2526

@@ -354,4 +355,49 @@ public String toString() {
354355
return "BINARY_AS_FLOAT16_IEEE_754_TOTAL_ORDER_COMPARATOR";
355356
}
356357
};
358+
359+
/**
360+
* Comparator for two timestamps encoded as INT96 (12-byte little-endian) binary.
361+
* Layout: first 8 bytes = nanoseconds within the day, last 4 bytes = Julian day.
362+
*
363+
* Two-level comparison, matching the INT96 timestamp sort order:
364+
* 1. Compare the last 4 bytes (Julian day) as a signed little-endian int32.
365+
* 2. If equal, compare the first 8 bytes (nanos) as a signed little-endian int64.
366+
*/
367+
static final PrimitiveComparator<Binary> BINARY_AS_INT96_TIMESTAMP_COMPARATOR = new BinaryComparator() {
368+
private static final long NANOSECONDS_PER_DAY = 86_400_000_000_000L;
369+
370+
@Override
371+
int compareBinary(Binary b1, Binary b2) {
372+
if (b1.length() != 12 || b2.length() != 12) {
373+
throw new IllegalArgumentException(
374+
"INT96 binary length must be 12, got " + b1.length() + " and " + b2.length());
375+
}
376+
377+
ByteBuffer bb1 = b1.toByteBuffer().slice();
378+
ByteBuffer bb2 = b2.toByteBuffer().slice();
379+
bb1.order(ByteOrder.LITTLE_ENDIAN);
380+
bb2.order(ByteOrder.LITTLE_ENDIAN);
381+
382+
int result = Integer.compare(bb1.getInt(8), bb2.getInt(8));
383+
if (result != 0) return result;
384+
385+
long nanos1 = bb1.getLong(0);
386+
long nanos2 = bb2.getLong(0);
387+
if (nanos1 < 0 || nanos1 > NANOSECONDS_PER_DAY) {
388+
throw new IllegalArgumentException(
389+
"Invalid nanos value (must be positive and less than 1 day): " + nanos1);
390+
}
391+
if (nanos2 < 0 || nanos2 > NANOSECONDS_PER_DAY) {
392+
throw new IllegalArgumentException(
393+
"Invalid nanos value (must be positive and less than 1 day): " + nanos2);
394+
}
395+
return Long.compare(nanos1, nanos2);
396+
}
397+
398+
@Override
399+
public String toString() {
400+
return "BINARY_AS_INT96_TIMESTAMP_COMPARATOR";
401+
}
402+
};
357403
}

parquet-column/src/main/java/org/apache/parquet/schema/PrimitiveType.java

Lines changed: 28 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -385,7 +385,9 @@ public <T, E extends Exception> T convert(PrimitiveTypeNameConverter<T, E> conve
385385

386386
@Override
387387
PrimitiveComparator<?> comparator(LogicalTypeAnnotation logicalType, ColumnOrder columnOrder) {
388-
return PrimitiveComparator.BINARY_AS_SIGNED_INTEGER_COMPARATOR;
388+
return columnOrder != null && columnOrder.getColumnOrderName() == ColumnOrderName.INT96_TIMESTAMP_ORDER
389+
? PrimitiveComparator.BINARY_AS_INT96_TIMESTAMP_COMPARATOR
390+
: PrimitiveComparator.BINARY_AS_SIGNED_INTEGER_COMPARATOR;
389391
}
390392
},
391393
FIXED_LEN_BYTE_ARRAY("getBinary", Binary.class) {
@@ -578,9 +580,14 @@ public PrimitiveType(
578580
this.decimalMeta = decimalMeta;
579581

580582
if (columnOrder == null) {
581-
columnOrder = primitive == PrimitiveTypeName.INT96 || originalType == OriginalType.INTERVAL
582-
? ColumnOrder.undefined()
583-
: ColumnOrder.typeDefined();
583+
if (primitive == PrimitiveTypeName.INT96) {
584+
// INT96 is only used for deprecated timestamps and supports no other semantics.
585+
columnOrder = ColumnOrder.int96TimestampOrder();
586+
} else if (originalType == OriginalType.INTERVAL) {
587+
columnOrder = ColumnOrder.undefined();
588+
} else {
589+
columnOrder = ColumnOrder.typeDefined();
590+
}
584591
} else if (columnOrder.getColumnOrderName() == ColumnOrderName.IEEE_754_TOTAL_ORDER) {
585592
Preconditions.checkArgument(
586593
primitive == PrimitiveTypeName.FLOAT || primitive == PrimitiveTypeName.DOUBLE,
@@ -629,10 +636,16 @@ public PrimitiveType(
629636
}
630637

631638
if (columnOrder == null) {
632-
columnOrder = primitive == PrimitiveTypeName.INT96
633-
|| logicalTypeAnnotation instanceof LogicalTypeAnnotation.IntervalLogicalTypeAnnotation
634-
? ColumnOrder.undefined()
635-
: ColumnOrder.typeDefined();
639+
if (primitive == PrimitiveTypeName.INT96) {
640+
// A plain INT96 is the legacy timestamp encoding; default it to the chronological order.
641+
// An annotated INT96 carries other semantics, so leave its order undefined.
642+
columnOrder =
643+
logicalTypeAnnotation == null ? ColumnOrder.int96TimestampOrder() : ColumnOrder.undefined();
644+
} else if (logicalTypeAnnotation instanceof LogicalTypeAnnotation.IntervalLogicalTypeAnnotation) {
645+
columnOrder = ColumnOrder.undefined();
646+
} else {
647+
columnOrder = ColumnOrder.typeDefined();
648+
}
636649
} else if (columnOrder.getColumnOrderName() == ColumnOrderName.IEEE_754_TOTAL_ORDER) {
637650
Preconditions.checkArgument(
638651
primitive == PrimitiveTypeName.FLOAT
@@ -651,9 +664,15 @@ public PrimitiveType(
651664
private ColumnOrder requireValidColumnOrder(ColumnOrder columnOrder) {
652665
if (primitive == PrimitiveTypeName.INT96) {
653666
Preconditions.checkArgument(
654-
columnOrder.getColumnOrderName() == ColumnOrderName.UNDEFINED,
667+
columnOrder.getColumnOrderName() == ColumnOrderName.UNDEFINED
668+
|| columnOrder.getColumnOrderName() == ColumnOrderName.INT96_TIMESTAMP_ORDER,
655669
"The column order %s is not supported by INT96",
656670
columnOrder);
671+
} else {
672+
Preconditions.checkArgument(
673+
columnOrder.getColumnOrderName() != ColumnOrderName.INT96_TIMESTAMP_ORDER,
674+
"The column order %s is only supported by INT96",
675+
columnOrder);
657676
}
658677
if (getLogicalTypeAnnotation() != null) {
659678
Preconditions.checkArgument(

parquet-column/src/test/java/org/apache/parquet/internal/column/columnindex/TestBinaryTruncator.java

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,18 @@ public void testContractNonStringTypes() {
8888
testTruncator(
8989
Types.required(FIXED_LEN_BYTE_ARRAY).length(12).as(INTERVAL).named("test_fixed_interval"), false);
9090
testTruncator(Types.required(BINARY).as(DECIMAL).precision(10).scale(2).named("test_binary_decimal"), false);
91-
testTruncator(Types.required(INT96).named("test_int96"), false);
91+
}
92+
93+
@Test
94+
public void testInt96() {
95+
// INT96 has a fixed 12-byte width and a chronological comparator (so it is excluded from the
96+
// variable-length checks above, like FLOAT16). Its truncator is a no-op: verify it returns the
97+
// value unchanged regardless of the requested length.
98+
BinaryTruncator int96Truncator =
99+
BinaryTruncator.getTruncator(Types.required(INT96).named("test_int96"));
100+
Binary int96Value = Binary.fromConstantByteArray(new byte[] {0, 0, 0, 0, 0, 0, 0, 0, 1, 2, 3, 4});
101+
assertThat(int96Truncator.truncateMin(int96Value, 4)).isSameAs(int96Value);
102+
assertThat(int96Truncator.truncateMax(int96Value, 4)).isSameAs(int96Value);
92103
}
93104

94105
@Test

parquet-column/src/test/java/org/apache/parquet/schema/TestPrimitiveComparator.java

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
import static org.apache.parquet.schema.PrimitiveComparator.BINARY_AS_FLOAT16_COMPARATOR;
2222
import static org.apache.parquet.schema.PrimitiveComparator.BINARY_AS_FLOAT16_IEEE_754_TOTAL_ORDER_COMPARATOR;
23+
import static org.apache.parquet.schema.PrimitiveComparator.BINARY_AS_INT96_TIMESTAMP_COMPARATOR;
2324
import static org.apache.parquet.schema.PrimitiveComparator.BINARY_AS_SIGNED_INTEGER_COMPARATOR;
2425
import static org.apache.parquet.schema.PrimitiveComparator.BOOLEAN_COMPARATOR;
2526
import static org.apache.parquet.schema.PrimitiveComparator.DOUBLE_COMPARATOR;
@@ -39,6 +40,8 @@
3940
import java.util.ArrayList;
4041
import java.util.Comparator;
4142
import java.util.List;
43+
import org.apache.parquet.TestUtils;
44+
import org.apache.parquet.example.data.simple.NanoTime;
4245
import org.apache.parquet.io.api.Binary;
4346
import org.junit.jupiter.api.Test;
4447

@@ -353,6 +356,48 @@ public void testBinaryAsSignedIntegerComparatorWithEquals() {
353356
}
354357
}
355358

359+
private static Binary int96(int julianDay, long nanosOfDay) {
360+
return new NanoTime(julianDay, nanosOfDay).toBinary();
361+
}
362+
363+
@Test
364+
public void testInt96TimestampComparator() {
365+
Binary[] valuesInAscendingOrder = {
366+
int96(Integer.MIN_VALUE, 0), // most negative julian day
367+
int96(-1, 86_399_999_999_999L), // negative julian days sort before day 0
368+
int96(0, 0), // start of the julian period
369+
int96(0, 86_399_999_999_999L), // same day, later time of day
370+
int96(2440000, 123L), // 1968-05-23T00:00:00.000000123, pre-epoch but positive julian day
371+
int96(2458850, 43_200_000_000_000L), // 2020-01-01T12:00:00
372+
int96(2458881, 39_600_000_000_000L), // 2020-02-01T11:00:00, later day even though earlier time of day
373+
int96(2458881, 39_600_000_000_001L), // 2020-02-01T11:00:00.000000001, nanos tie-break
374+
int96(Integer.MAX_VALUE, 86_399_999_999_999L)
375+
};
376+
377+
for (int i = 0; i < valuesInAscendingOrder.length; ++i) {
378+
for (int j = 0; j < valuesInAscendingOrder.length; ++j) {
379+
assertThat(Integer.signum(BINARY_AS_INT96_TIMESTAMP_COMPARATOR.compare(
380+
valuesInAscendingOrder[i], valuesInAscendingOrder[j])))
381+
.as("comparing value " + i + " to value " + j)
382+
.isEqualTo(Integer.signum(Integer.compare(i, j)));
383+
}
384+
}
385+
}
386+
387+
@Test
388+
public void testInt96TimestampComparatorRejectsInvalidNanos() {
389+
// Same Julian day so the comparator reaches the nanos validation instead of
390+
// returning early on the day comparison.
391+
Binary valid = int96(0, 0);
392+
for (long invalidNanos : new long[] {-1L, Long.MIN_VALUE, 86_400_000_000_001L, Long.MAX_VALUE}) {
393+
Binary invalid = int96(0, invalidNanos);
394+
TestUtils.assertThrows(
395+
"Expected IllegalArgumentException for nanos=" + invalidNanos,
396+
IllegalArgumentException.class,
397+
() -> BINARY_AS_INT96_TIMESTAMP_COMPARATOR.compare(valid, invalid));
398+
}
399+
}
400+
356401
@Test
357402
public void testFloat16Comparator() {
358403
Binary[] valuesInAscendingOrder = {

parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java

Lines changed: 28 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,7 @@
8888
import org.apache.parquet.format.GeometryType;
8989
import org.apache.parquet.format.GeospatialStatistics;
9090
import org.apache.parquet.format.IEEE754TotalOrder;
91+
import org.apache.parquet.format.Int96TimestampOrder;
9192
import org.apache.parquet.format.IntType;
9293
import org.apache.parquet.format.KeyValue;
9394
import org.apache.parquet.format.LogicalType;
@@ -146,6 +147,7 @@ public class ParquetMetadataConverter {
146147

147148
private static final TypeDefinedOrder TYPE_DEFINED_ORDER = new TypeDefinedOrder();
148149
private static final IEEE754TotalOrder IEEE_754_TOTAL_ORDER = new IEEE754TotalOrder();
150+
private static final Int96TimestampOrder INT96_TIMESTAMP_ORDER = new Int96TimestampOrder();
149151
public static final MetadataFilter NO_FILTER = new NoFilter();
150152
public static final MetadataFilter SKIP_ROW_GROUPS = new SkipMetadataFilter();
151153
public static final long MAX_STATS_SIZE = 4096; // limit stats to 4k
@@ -290,6 +292,9 @@ private List<ColumnOrder> getColumnOrders(MessageType schema) {
290292
case IEEE_754_TOTAL_ORDER:
291293
columnOrder.setIEEE_754_TOTAL_ORDER(IEEE_754_TOTAL_ORDER);
292294
break;
295+
case INT96_TIMESTAMP_ORDER:
296+
columnOrder.setINT96_TIMESTAMP_ORDER(INT96_TIMESTAMP_ORDER);
297+
break;
293298
case UNDEFINED:
294299
// Use TypeDefinedOrder if some types (e.g. INT96) have undefined column orders.
295300
columnOrder.setTYPE_ORDER(TYPE_DEFINED_ORDER);
@@ -911,8 +916,16 @@ private static byte[] tuncateMax(BinaryTruncator truncator, int truncateLength,
911916
}
912917

913918
private static boolean isMinMaxStatsSupported(PrimitiveType type) {
914-
return type.columnOrder().getColumnOrderName() == ColumnOrderName.TYPE_DEFINED_ORDER
915-
|| type.columnOrder().getColumnOrderName() == ColumnOrderName.IEEE_754_TOTAL_ORDER;
919+
switch (type.columnOrder().getColumnOrderName()) {
920+
case TYPE_DEFINED_ORDER:
921+
case IEEE_754_TOTAL_ORDER:
922+
case INT96_TIMESTAMP_ORDER:
923+
return true;
924+
case UNDEFINED:
925+
return false;
926+
default:
927+
throw new IllegalArgumentException("Unknown column order: " + type.columnOrder());
928+
}
916929
}
917930

918931
/**
@@ -2062,7 +2075,17 @@ private void buildChildren(
20622075
|| schemaElement.converted_type == ConvertedType.INTERVAL)) {
20632076
columnOrder = org.apache.parquet.schema.ColumnOrder.undefined();
20642077
}
2078+
// INT96_TIMESTAMP_ORDER is only valid for INT96 columns, ignore it anywhere else.
2079+
if (columnOrder.getColumnOrderName() == ColumnOrderName.INT96_TIMESTAMP_ORDER
2080+
&& schemaElement.type != Type.INT96) {
2081+
columnOrder = org.apache.parquet.schema.ColumnOrder.undefined();
2082+
}
20652083
primitiveBuilder.columnOrder(columnOrder);
2084+
} else if (schemaElement.type == Type.INT96) {
2085+
// A footer without column orders predates INT96_TIMESTAMP_ORDER, so an INT96 column here
2086+
// must not inherit the (chronological) construction-time default: its stats, if any, were
2087+
// written under the legacy order and must be ignored.
2088+
primitiveBuilder.columnOrder(org.apache.parquet.schema.ColumnOrder.undefined());
20662089
}
20672090
childBuilder = primitiveBuilder;
20682091
} else {
@@ -2118,6 +2141,9 @@ private static org.apache.parquet.schema.ColumnOrder fromParquetColumnOrder(Colu
21182141
if (columnOrder.isSetIEEE_754_TOTAL_ORDER()) {
21192142
return org.apache.parquet.schema.ColumnOrder.ieee754TotalOrder();
21202143
}
2144+
if (columnOrder.isSetINT96_TIMESTAMP_ORDER()) {
2145+
return org.apache.parquet.schema.ColumnOrder.int96TimestampOrder();
2146+
}
21212147
// The column order is not yet supported by this API
21222148
return org.apache.parquet.schema.ColumnOrder.undefined();
21232149
}

0 commit comments

Comments
 (0)