diff --git a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/RecordConverter.java b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/RecordConverter.java index ab3d5aa9bb43..516f5f49abfa 100644 --- a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/RecordConverter.java +++ b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/RecordConverter.java @@ -95,11 +95,19 @@ class RecordConverter { private final NameMapping nameMapping; private final IcebergSinkConfig config; private final Map> structNameMap = Maps.newHashMap(); + // Parquet stores UUIDs as a 16-byte fixed; other formats keep the UUID logical type. The write + // file format is fixed for the converter's lifetime, so resolve this once instead of per value. + private final boolean writeUuidAsBytes; RecordConverter(Table table, IcebergSinkConfig config) { this.tableSchema = table.schema(); this.nameMapping = createNameMapping(table); this.config = config; + this.writeUuidAsBytes = + FileFormat.PARQUET + .name() + .toLowerCase(Locale.ROOT) + .equals(config.writeProps().get(TableProperties.DEFAULT_FILE_FORMAT)); } Record convert(Object data) { @@ -420,10 +428,7 @@ protected Object convertUUID(Object value) { throw new IllegalArgumentException("Cannot convert to UUID: " + value.getClass().getName()); } - if (FileFormat.PARQUET - .name() - .toLowerCase(Locale.ROOT) - .equals(config.writeProps().get(TableProperties.DEFAULT_FILE_FORMAT))) { + if (writeUuidAsBytes) { return UUIDUtil.convert(uuid); } else { return uuid; diff --git a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestRecordConverter.java b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestRecordConverter.java index 9b91ba61c167..b8a46274f433 100644 --- a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestRecordConverter.java +++ b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestRecordConverter.java @@ -220,6 +220,8 @@ public static void beforeAll() { public void before() { this.config = mock(IcebergSinkConfig.class); when(config.jsonConverter()).thenReturn(JSON_CONVERTER); + // production writeProps() is never null; default it so RecordConverter construction succeeds + when(config.writeProps()).thenReturn(ImmutableMap.of()); } @Test