This is an automated email from the ASF dual-hosted git repository.
snuyanzin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 30a02eaa93c [FLINK-40515][core] Extend BinaryVariant to support UUID
30a02eaa93c is described below
commit 30a02eaa93c99ec3fcd179b67560e6bd8504d334
Author: Moritz Manner <[email protected]>
AuthorDate: Thu Sep 10 15:03:37 2026 +0200
[FLINK-40515][core] Extend BinaryVariant to support UUID
---
.../apache/flink/types/variant/BinaryVariant.java | 12 ++++
.../flink/types/variant/BinaryVariantBuilder.java | 8 +++
.../variant/BinaryVariantInternalBuilder.java | 13 +++++
.../flink/types/variant/BinaryVariantUtil.java | 47 ++++++++++++++++
.../org/apache/flink/types/variant/Variant.java | 12 +++-
.../apache/flink/types/variant/VariantBuilder.java | 4 ++
.../flink/types/variant/BinaryVariantTest.java | 64 ++++++++++++++++++++++
7 files changed, 159 insertions(+), 1 deletion(-)
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
index 6cccd45d524..e30217e9ea5 100644
--- a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
+++ b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
@@ -38,6 +38,7 @@ import java.util.Arrays;
import java.util.Base64;
import java.util.List;
import java.util.Objects;
+import java.util.UUID;
import static
org.apache.flink.types.variant.BinaryVariantUtil.BINARY_SEARCH_THRESHOLD;
import static org.apache.flink.types.variant.BinaryVariantUtil.SIZE_LIMIT;
@@ -224,6 +225,12 @@ public final class BinaryVariant implements Variant {
return BinaryVariantUtil.getBinary(value, pos);
}
+ @Override
+ public UUID getUuid() throws VariantTypeException {
+ checkType(Type.UUID, getType());
+ return BinaryVariantUtil.getUuid(value, pos);
+ }
+
@Override
public Object get() throws VariantTypeException {
switch (getType()) {
@@ -259,6 +266,8 @@ public final class BinaryVariant implements Variant {
return getInstant();
case BYTES:
return getBytes();
+ case UUID:
+ return getUuid();
default:
throw new VariantTypeException(
String.format("Expecting a primitive variant but got
%s", getType()));
@@ -459,6 +468,9 @@ public final class BinaryVariant implements Variant {
Base64.getEncoder()
.encodeToString(BinaryVariantUtil.getBinary(value, pos)));
break;
+ case UUID:
+ appendQuoted(sb, BinaryVariantUtil.getUuid(value,
pos).toString());
+ break;
default:
throw unexpectedType(BinaryVariantUtil.getType(value, pos));
}
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantBuilder.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantBuilder.java
index c2b9ef59faf..c467c9adad2 100644
---
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantBuilder.java
+++
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantBuilder.java
@@ -29,6 +29,7 @@ import java.time.LocalTime;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
+import java.util.UUID;
/** Builder for binary encoded variant. */
@Internal
@@ -168,6 +169,13 @@ public class BinaryVariantBuilder implements
VariantBuilder {
return builder.build();
}
+ @Override
+ public Variant of(UUID uuid) {
+ BinaryVariantInternalBuilder builder = new
BinaryVariantInternalBuilder(false);
+ builder.appendUuid(uuid);
+ return builder.build();
+ }
+
@Override
public Variant ofNull() {
BinaryVariantInternalBuilder builder = new
BinaryVariantInternalBuilder(false);
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java
index 5b0cadddc3c..a61763a93b6 100644
---
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java
+++
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java
@@ -35,6 +35,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
+import java.util.UUID;
import static org.apache.flink.types.variant.BinaryVariantUtil.ARRAY;
import static org.apache.flink.types.variant.BinaryVariantUtil.BASIC_TYPE_MASK;
@@ -328,6 +329,18 @@ public class BinaryVariantInternalBuilder {
writePos += binary.length;
}
+ public void appendUuid(UUID uuid) {
+ checkCapacity(1 + 16);
+ writeBuffer[writePos++] = primitiveHeader(BinaryVariantUtil.UUID);
+ // The variant spec stores UUIDs as 16 big-endian bytes: the most
significant 8 bytes
+ // followed by the least significant 8 bytes. UUID is the only
primitive that is not
+ // little-endian, so we use the big-endian writer instead of writeLong.
+ BinaryVariantUtil.writeLongBigEndian(writeBuffer, writePos,
uuid.getMostSignificantBits());
+ BinaryVariantUtil.writeLongBigEndian(
+ writeBuffer, writePos + 8, uuid.getLeastSignificantBits());
+ writePos += 16;
+ }
+
// Add a key to the variant dictionary. If the key already exists, the
dictionary is not
// modified.
// In either case, return the id of the key.
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
index 8bf7ef61fab..575f54f2550 100644
---
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
+++
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
@@ -28,6 +28,7 @@ import java.time.format.DateTimeFormatter;
import java.time.format.DateTimeFormatterBuilder;
import java.util.Arrays;
import java.util.Locale;
+import java.util.UUID;
/* This file is based on source code from the Spark Project
(http://spark.apache.org/), licensed by the Apache
* Software Foundation (ASF) under the Apache License, Version 2.0. See the
NOTICE file distributed with this work for
@@ -201,6 +202,9 @@ public class BinaryVariantUtil {
*/
public static final int TIMESTAMP_NS = 19;
+ /** UUID value */
+ public static final int UUID = 20;
+
public static final byte VERSION = 1;
/** The lower 4 bits of the first metadata byte contain the version. */
@@ -249,6 +253,16 @@ public class BinaryVariantUtil {
}
}
+ /**
+ * Write a big-endian 8-byte long value to {@code bytes[pos, pos + 8)}.
Variant primitives are
+ * little-endian except for UUID, which the spec stores as 16 big-endian
bytes.
+ */
+ static void writeLongBigEndian(byte[] bytes, int pos, long value) {
+ for (int i = 0; i < 8; ++i) {
+ bytes[pos + i] = (byte) ((value >>> (8 * (7 - i))) & 0xFF);
+ }
+ }
+
public static byte primitiveHeader(int type) {
return (byte) (type << 2 | PRIMITIVE);
}
@@ -322,6 +336,18 @@ public class BinaryVariantUtil {
return result;
}
+ /**
+ * Read a big-endian 8-byte long value from {@code bytes[pos, pos + 8)}.
The caller must ensure
+ * those 8 bytes are within bounds.
+ */
+ private static long readLongBigEndian(byte[] bytes, int pos) {
+ long result = 0;
+ for (int i = 0; i < 8; ++i) {
+ result = (result << 8) | (bytes[pos + i] & 0xFF);
+ }
+ return result;
+ }
+
/**
* Read a little-endian unsigned int value from {@code bytes[pos, pos +
numBytes)}. The value
* must fit into a non-negative int ({@code [0, Integer.MAX_VALUE]}).
@@ -403,6 +429,8 @@ public class BinaryVariantUtil {
return Type.TIMESTAMP_LTZ_NS;
case TIMESTAMP_NS:
return Type.TIMESTAMP_NS;
+ case UUID:
+ return Type.UUID;
default:
throw unknownPrimitiveTypeInVariant(typeInfo);
}
@@ -470,6 +498,8 @@ public class BinaryVariantUtil {
return 6;
case DECIMAL8:
return 10;
+ case UUID:
+ return 17;
case DECIMAL16:
return 18;
case BINARY:
@@ -674,6 +704,23 @@ public class BinaryVariantUtil {
throw unexpectedType(Type.STRING);
}
+ public static UUID getUuid(byte[] value, int pos) {
+ // The header byte and the 16 UUID bytes occupy value[pos, pos + 16],
so checking the first
+ // and last index once covers the whole read.
+ checkIndex(pos, value.length);
+ checkIndex(pos + 16, value.length);
+ int basicType = value[pos] & BASIC_TYPE_MASK;
+ int typeInfo = (value[pos] >> BASIC_TYPE_BITS) & TYPE_INFO_MASK;
+ if (basicType != PRIMITIVE || typeInfo != UUID) {
+ throw unexpectedType(Type.UUID);
+ }
+ // The 16 UUID bytes are big-endian (most significant bytes first):
the most significant 8
+ // bytes followed by the least significant 8 bytes.
+ long msb = readLongBigEndian(value, pos + 1);
+ long lsb = readLongBigEndian(value, pos + 9);
+ return new UUID(msb, lsb);
+ }
+
/** A handler that receives the decoded header fields of a variant object.
*/
public interface ObjectHandler<T> {
/**
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
b/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
index cad26314ad3..6a01b75224f 100644
--- a/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
+++ b/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
@@ -27,6 +27,7 @@ import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
import java.util.List;
+import java.util.UUID;
/**
* Variant represent a semi-structured data.
@@ -171,6 +172,14 @@ public interface Variant extends Serializable {
*/
byte[] getBytes() throws VariantTypeException;
+ /**
+ * Get the scalar value of variant as UUID, if the variant type is {@link
Type#UUID}.
+ *
+ * @throws VariantTypeException If this variant is not a scalar value or
is not {@link
+ * Type#UUID}.
+ */
+ UUID getUuid() throws VariantTypeException;
+
/**
* Get the scalar value of variant.
*
@@ -249,7 +258,8 @@ public interface Variant extends Serializable {
TIMESTAMP_LTZ,
TIMESTAMP_NS,
TIMESTAMP_LTZ_NS,
- BYTES
+ BYTES,
+ UUID
}
static VariantBuilder newBuilder() {
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/VariantBuilder.java
b/flink-core/src/main/java/org/apache/flink/types/variant/VariantBuilder.java
index d73d0fc3c25..abda3b8491d 100644
---
a/flink-core/src/main/java/org/apache/flink/types/variant/VariantBuilder.java
+++
b/flink-core/src/main/java/org/apache/flink/types/variant/VariantBuilder.java
@@ -25,6 +25,7 @@ import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
+import java.util.UUID;
/** Builder for variants. */
@PublicEvolving
@@ -72,6 +73,9 @@ public interface VariantBuilder {
/** Create a variant from a LocalTime. Sub-microsecond precision is
truncated. */
Variant of(LocalTime localTime);
+ /** Create a variant from a UUID. */
+ Variant of(UUID uuid);
+
/** Create a variant of null. */
Variant ofNull();
diff --git
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
index 76c5222c168..5c4a6453dcc 100644
---
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
+++
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
@@ -34,6 +34,7 @@ import java.time.LocalTime;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
import java.util.Collections;
+import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -105,6 +106,10 @@ class BinaryVariantTest {
assertThat(builder.of(localTime).getTime()).isEqualTo(localTime);
assertThat(builder.of(localTime).get()).isEqualTo(localTime);
+ UUID uuid = UUID.randomUUID();
+ assertThat(builder.of(uuid).getUuid()).isEqualTo(uuid);
+ assertThat(builder.of(uuid).get()).isEqualTo(uuid);
+
assertThat(builder.ofNull().get()).isEqualTo(null);
assertThat(builder.ofNull().isNull()).isTrue();
}
@@ -289,6 +294,7 @@ class BinaryVariantTest {
LocalTime localTime = LocalTime.of(13, 45, 30, 123456789);
Instant nanoInstant = Instant.EPOCH.plusNanos(123456789);
LocalDateTime nanoLocalDateTime = LocalDateTime.of(2000, 1, 1, 0, 0,
0, 123456789);
+ UUID uuid = UUID.fromString("00112233-4455-6677-8899-aabbccddeeff");
assertThat(builder.of((byte) 1).toJson()).isEqualTo("1");
assertThat(builder.of((short) 1).toJson()).isEqualTo("1");
@@ -308,9 +314,46 @@ class BinaryVariantTest {
assertThat(builder.of(nanoLocalDateTime).toJson())
.isEqualTo("\"2000-01-01T00:00:00.123456789\"");
assertThat(builder.of("hello".getBytes()).toJson()).isEqualTo("\"aGVsbG8=\"");
+
assertThat(builder.of(uuid).toJson()).isEqualTo("\"00112233-4455-6677-8899-aabbccddeeff\"");
assertThat(builder.ofNull().toJson()).isEqualTo("null");
}
+ @Test
+ void testUuidDecodeFromSpecBytes() {
+ // Interop check against the shared variant wire format: this is the
exact byte sequence
+ // from Iceberg's TestSerializedPrimitives#testUUID (primitive header
for type 20 followed
+ // by 16 big-endian UUID bytes). Decoding it must produce the same
UUID, which proves Flink
+ // reads variants written by other implementations of the spec.
+ //
https://github.com/apache/iceberg/blob/9da109dd2537e77e8e5034068933575aa4e235ff/api/src/test/java/org/apache/iceberg/variants/TestSerializedPrimitives.java#L586
+ byte[] value = {
+ BinaryVariantUtil.primitiveHeader(BinaryVariantUtil.UUID),
+ (byte) 0xf2,
+ 0x4f,
+ (byte) 0x9b,
+ 0x64,
+ (byte) 0x81,
+ (byte) 0xfa,
+ 0x49,
+ (byte) 0xd1,
+ (byte) 0xb7,
+ 0x4e,
+ (byte) 0x8c,
+ 0x09,
+ (byte) 0xa6,
+ (byte) 0xe3,
+ 0x1c,
+ 0x56
+ };
+ // A primitive carries no dictionary keys, so reuse an empty metadata
block.
+ byte[] metadata = ((BinaryVariant) builder.of(0)).getMetadata();
+ Variant variant = new BinaryVariant(value, metadata);
+
+ UUID expected =
UUID.fromString("f24f9b64-81fa-49d1-b74e-8c09a6e31c56");
+ assertThat(variant.getType()).isEqualTo(Variant.Type.UUID);
+ assertThat(variant.getUuid()).isEqualTo(expected);
+ assertThat(variant.get()).isEqualTo(expected);
+ }
+
@Test
void testToJsonNested() {
Variant variant =
@@ -426,6 +469,27 @@ class BinaryVariantTest {
.hasMessage("Expected type DOUBLE but got FLOAT");
}
+ @Test
+ void testUuidGetThrowException() {
+ // Reading a UUID from a non-UUID variant, and reading another type
from a UUID variant,
+ // must both fail with a type exception.
+ assertThatThrownBy(builder.of(10)::getUuid)
+ .isInstanceOf(VariantTypeException.class)
+ .hasMessage("Expected type UUID but got INT");
+
+ assertThatThrownBy(builder.of(UUID.randomUUID())::getString)
+ .isInstanceOf(VariantTypeException.class)
+ .hasMessage("Expected type STRING but got UUID");
+
+ // A UUID header followed by fewer than 16 bytes is malformed and must
be rejected
+ byte[] truncated = {
+ BinaryVariantUtil.primitiveHeader(BinaryVariantUtil.UUID), 0x00,
0x01, 0x02, 0x03
+ };
+ assertThatThrownBy(() -> BinaryVariantUtil.getUuid(truncated, 0))
+ .isInstanceOf(VariantTypeException.class)
+ .hasMessage("MALFORMED_VARIANT");
+ }
+
@Test
void testJavaSerialization() throws Exception {
Variant variant =