This is an automated email from the ASF dual-hosted git repository.
laskoviymishka pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new 43c5436762 Kafka Connect: Fix invalid decimal type inferred for some
BigDecimal values (#16606)
43c5436762 is described below
commit 43c54367627fd1d59fbda03010318af75c711ea2
Author: Vova Kolmakov <[email protected]>
AuthorDate: Mon Aug 3 22:55:55 2026 +0700
Kafka Connect: Fix invalid decimal type inferred for some BigDecimal values
(#16606)
* Kafka Connect: Fix invalid decimal type inferred for some BigDecimal
values
Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
* Kafka Connect: Bound the precision of an inferred decimal type
Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Generated-by: Claude Code (claude-opus-5)
---------
Co-authored-by: Vova Kolmakov <[email protected]>
Co-authored-by: Claude Opus 4.8 (1M context) <[email protected]>
---
.../apache/iceberg/connect/data/SchemaUtils.java | 30 ++++++++++++++--
.../iceberg/connect/data/TestSchemaUtils.java | 37 +++++++++++++++++++
.../iceberg/connect/data/TestSinkWriter.java | 42 ++++++++++++++++++++++
3 files changed, 107 insertions(+), 2 deletions(-)
diff --git
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/SchemaUtils.java
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/SchemaUtils.java
index 89d7878172..5bcaed64ac 100644
---
a/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/SchemaUtils.java
+++
b/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/data/SchemaUtils.java
@@ -72,6 +72,8 @@ class SchemaUtils {
private static final Pattern TRANSFORM_REGEX =
Pattern.compile("(\\w+)\\((.+)\\)");
+ private static final int MAX_DECIMAL_PRECISION = 38;
+
static PrimitiveType needsDataTypeUpdate(Type currentIcebergType, Schema
valueSchema) {
if (currentIcebergType.typeId() == TypeID.FLOAT && valueSchema.type() ==
Schema.Type.FLOAT64) {
return DoubleType.get();
@@ -305,8 +307,7 @@ class SchemaUtils {
} else if (value instanceof Boolean) {
return BooleanType.get();
} else if (value instanceof BigDecimal) {
- BigDecimal bigDecimal = (BigDecimal) value;
- return DecimalType.of(bigDecimal.precision(), bigDecimal.scale());
+ return inferDecimalType((BigDecimal) value);
} else if (value instanceof Integer || value instanceof Long) {
return LongType.get();
} else if (value instanceof Float || value instanceof Double) {
@@ -349,6 +350,31 @@ class SchemaUtils {
}
}
+ /**
+ * BigDecimal does not satisfy Iceberg's 0 <= scale <= precision
<= 38 invariant: a
+ * value < 1 has a precision smaller than its scale ("0.001" is
precision 1, scale 3), and an
+ * exponential value has a negative scale ("1E+2" is scale -2). Both are
normalized here, the
+ * same way Spark normalizes a BigDecimal in Decimal.set. A value that
needs more than 38 digits
+ * cannot be represented, so its type is reported as unknown.
+ */
+ private static Type inferDecimalType(BigDecimal value) {
+ // widened to long because the subtraction below overflows int for a
pathological scale,
+ // e.g. new BigDecimal(BigInteger.ONE, Integer.MIN_VALUE)
+ long scale = value.scale();
+ long precision = value.precision();
+ if (scale < 0) {
+ precision -= scale;
+ scale = 0;
+ }
+
+ precision = Math.max(precision, scale);
+ if (precision > MAX_DECIMAL_PRECISION) {
+ return null;
+ }
+
+ return DecimalType.of((int) precision, (int) scale);
+ }
+
private int nextId() {
return fieldId++;
}
diff --git
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSchemaUtils.java
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSchemaUtils.java
index 9443ed4676..602228c481 100644
---
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSchemaUtils.java
+++
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSchemaUtils.java
@@ -30,6 +30,7 @@ import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import java.math.BigDecimal;
+import java.math.BigInteger;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
@@ -340,4 +341,40 @@ public class TestSchemaUtils {
Schema uuidSchema = SchemaBuilder.string().name("uuid").build();
assertThat(SchemaUtils.toIcebergType(uuidSchema,
config)).isInstanceOf(UUIDType.class);
}
+
+ @Test
+ public void testInferIcebergTypeSmallDecimal() {
+ IcebergSinkConfig config = mock(IcebergSinkConfig.class);
+
+ // BigDecimal("0.001") has precision 1, smaller than its scale 3;
+ // Iceberg requires scale <= precision, so precision is widened to the
scale
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("0.001"), config))
+ .isEqualTo(DecimalType.of(3, 3));
+
+ // BigDecimal("1E+2") has a negative scale (-2); normalized to scale 0,
+ // the same decimal(3, 0) as new BigDecimal("100")
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("1E+2"), config))
+ .isEqualTo(DecimalType.of(3, 0));
+ }
+
+ @Test
+ public void testInferIcebergTypeDecimalOutOfRange() {
+ IcebergSinkConfig config = mock(IcebergSinkConfig.class);
+
+ // the widest values that can still be represented
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("1E+37"), config))
+ .isEqualTo(DecimalType.of(38, 0));
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("1E-38"), config))
+ .isEqualTo(DecimalType.of(38, 38));
+
+ // one digit past the limit on either branch: no type can be inferred
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("1E+38"),
config)).isNull();
+ assertThat(SchemaUtils.inferIcebergType(new BigDecimal("1E-39"),
config)).isNull();
+
+ // a scale of Integer.MIN_VALUE must not overflow the precision
normalization. it is built
+ // here rather than parsed from "1E+2147483648", which only yields this
scale on Java 21+
+ // (Java 17 rejects that exponent with a NumberFormatException)
+ BigDecimal minScale = new BigDecimal(BigInteger.ONE, Integer.MIN_VALUE);
+ assertThat(SchemaUtils.inferIcebergType(minScale, config)).isNull();
+ }
}
diff --git
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSinkWriter.java
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSinkWriter.java
index 09f7a373d5..97eb912e89 100644
---
a/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSinkWriter.java
+++
b/kafka-connect/kafka-connect/src/test/java/org/apache/iceberg/connect/data/TestSinkWriter.java
@@ -26,6 +26,7 @@ import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import java.io.IOException;
+import java.math.BigDecimal;
import java.time.Instant;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
@@ -235,6 +236,47 @@ public class TestSinkWriter {
assertThat(writerResults).hasSize(0);
}
+ @Test
+ public void testEvolveAddsFractionalDecimalColumn() {
+ IcebergSinkConfig config = mock(IcebergSinkConfig.class);
+ when(config.tableConfig(any())).thenReturn(mock(TableSinkConfig.class));
+
when(config.tables()).thenReturn(ImmutableList.of(TABLE_IDENTIFIER.toString()));
+ when(config.evolveSchemaEnabled()).thenReturn(true);
+
+ // a new column whose value is a fractional decimal < 1:
BigDecimal("0.001") has precision 1
+ // and scale 3. DecimalType.of validates only precision <= 38, so before
the fix the column
+ // evolved to a malformed decimal(1, 3) without error and the write below
failed in the
+ // Parquet writer; after the fix it evolves to decimal(3, 3) and the
record is written.
+ Map<String, Object> value = ImmutableMap.of("amount", new
BigDecimal("0.001"));
+
+ List<IcebergWriterResult> writerResults = sinkWriterTest(value, config);
+ assertThat(writerResults).isNotEmpty();
+
+ // the column evolved to a valid decimal that can hold 0.001 (scale <=
precision)
+ Types.NestedField added =
catalog.loadTable(TABLE_IDENTIFIER).schema().findField("amount");
+ assertThat(added).isNotNull();
+ assertThat(added.type()).isEqualTo(Types.DecimalType.of(3, 3));
+ }
+
+ @Test
+ public void testEvolveAddsExponentialDecimalColumn() {
+ IcebergSinkConfig config = mock(IcebergSinkConfig.class);
+ when(config.tableConfig(any())).thenReturn(mock(TableSinkConfig.class));
+
when(config.tables()).thenReturn(ImmutableList.of(TABLE_IDENTIFIER.toString()));
+ when(config.evolveSchemaEnabled()).thenReturn(true);
+
+ // BigDecimal("1E+2") has a negative scale (-2), which is normalized to
decimal(3, 0), the
+ // same type inferred for new BigDecimal("100"); the value is rescaled
when it is written
+ Map<String, Object> value = ImmutableMap.of("amount", new
BigDecimal("1E+2"));
+
+ List<IcebergWriterResult> writerResults = sinkWriterTest(value, config);
+ assertThat(writerResults).isNotEmpty();
+
+ Types.NestedField added =
catalog.loadTable(TABLE_IDENTIFIER).schema().findField("amount");
+ assertThat(added).isNotNull();
+ assertThat(added.type()).isEqualTo(Types.DecimalType.of(3, 0));
+ }
+
private List<IcebergWriterResult> sinkWriterTest(
Map<String, Object> value, IcebergSinkConfig config) {
IcebergWriterResult writeResult =