This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11415-ff598706f248b03cfab93b8188b2f763bf23e053 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 455d3002831812819726d16376bc5ad588307bfd Author: Jast <[email protected]> AuthorDate: Sat Sep 19 00:59:11 2026 +0000 [Fix][Format] Handle heterogeneous numeric JSON fields (#11415) Co-authored-by: zhangshenghang <[email protected]> Co-authored-by: zhangshenghang <[email protected]> --- .../introduction/concepts/incompatible-changes.md | 7 +++ .../introduction/concepts/incompatible-changes.md | 7 +++ .../seatunnel/format/json/RowToJsonConverters.java | 68 +++++++++++++++++++--- .../format/json/JsonRowDataSerDeSchemaTest.java | 48 +++++++++++++++ 4 files changed, 123 insertions(+), 7 deletions(-) diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index f2fedd0aca..99181067e1 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -325,6 +325,13 @@ You need to check this document before you upgrade to related version. — or filter the offending rows out upstream. Queries that worked around the `ABS` / `SIGN` rejection by casting (`ABS(CAST(tiny_col AS INT))`) continue to work unchanged and can be simplified at your convenience. +### Format Changes + +- **Breaking Change: JSON serialization of numeric fields now follows the runtime value type** + - **Affected component**: `seatunnel-formats/seatunnel-format-json` (`RowToJsonConverters`) - affects every connector that serializes rows with the JSON format (for example Kafka, RabbitMQ, Pulsar, and file JSON sinks) + - **Description**: Previously, a field declared as a numeric type in the catalog (`TINYINT`, `SMALLINT`, `INT`, `BIGINT`, `FLOAT`, `DOUBLE`, `DECIMAL`) was serialized by blindly casting the runtime value to the Java type implied by the declared type (for example `(long) value` for `BIGINT`). In multi-table jobs (for example CDC jobs writing JSON to RabbitMQ/Kafka) where several tables share one catalog schema but carry different physical column types, a `String` or `BigDecimal` runtime [...] + - **Impact**: Heterogeneous numeric values that previously crashed the job with `ClassCastException` now serialize successfully, and the emitted JSON numeric shape follows the runtime value rather than the declared column type (a `String` or `BigDecimal` value in a `BIGINT` column keeps its exact numeric value). Runtime values that can neither be represented as a number nor parsed from text (for example `byte[]`, `Map`, `LocalDateTime`) now fail fast with a typed `SeaTunnelJsonFormatEx [...] + ### Engine Behavior Changes ### Dependency Upgrades diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index bec69ded47..8a3b55c330 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -287,6 +287,13 @@ `ROUND(CAST(tiny_col AS INT), -1)`——或者在上游过滤掉这些行。此前为绕开 `ABS` / `SIGN` 拒绝而使用的强制转换 (`ABS(CAST(tiny_col AS INT))`)仍然可以正常工作,可以在方便时再简化。 +### 格式变更 + +- **破坏性变更:JSON 数值字段的序列化改为按运行时实际类型处理** + - **影响范围**:`seatunnel-formats/seatunnel-format-json`(`RowToJsonConverters`)--影响所有以 JSON 格式序列化行的连接器(例如 Kafka、RabbitMQ、Pulsar 及文件 JSON Sink)。 + - **变更说明**:以前,目录 Schema 中声明为数值类型(`TINYINT`、`SMALLINT`、`INT`、`BIGINT`、`FLOAT`、`DOUBLE`、`DECIMAL`)的字段,序列化时会把运行时值强制转换为声明类型对应的 Java 类型(例如 `BIGINT` 直接 `(long) value`)。在多表作业(例如多表 CDC 作业写 JSON 到 RabbitMQ/Kafka)中,多张表共享同一份目录 Schema 但物理列类型不一致时,`String` 或 `BigDecimal` 运行时值会抛出原始 `ClassCastException` 并导致作业失败。现在数值字段按运行时实际类型序列化:任意数值包装类型(`Byte`、`Short`、`Integer`、`Long`、`Float`、`Double`、`BigInteger`、`BigDecimal`)输出为对应的 JSON 数字;可解析为数字的字符串会解析成 JSON 数字,无法解析的文本则输出为 JSON 字符串;声明为 `DECIMAL` 的字段遇到 `Float`/`Dou [...] + - **影响**:以前因 `ClassCastException` 崩溃的异构数值现在可以正常序列化,输出的 JSON 数值形态跟随运行时值而非声明的列类型(`BIGINT` 列中的 `String` 或 `BigDecimal` 值会保留其精确数值)。既不能表示为数字、也无法从文本解析的运行时值(例如 `byte[]`、`Map`、`LocalDateTime`)将以类型化的 `SeaTunnelJsonFormatException`(`UNSUPPORTED_DATA_TYPE`)快速失败,替代原来的原始 `ClassCastException`。假定 JSON 数值形态始终与声明列类型一致的下游消费方需要重新评估。(#11415) + ### 引擎行为变更 ### 依赖升级 diff --git a/seatunnel-formats/seatunnel-format-json/src/main/java/org/apache/seatunnel/format/json/RowToJsonConverters.java b/seatunnel-formats/seatunnel-format-json/src/main/java/org/apache/seatunnel/format/json/RowToJsonConverters.java index 5aaf0a4995..b943149f8b 100644 --- a/seatunnel-formats/seatunnel-format-json/src/main/java/org/apache/seatunnel/format/json/RowToJsonConverters.java +++ b/seatunnel-formats/seatunnel-format-json/src/main/java/org/apache/seatunnel/format/json/RowToJsonConverters.java @@ -34,6 +34,7 @@ import org.apache.seatunnel.format.json.exception.SeaTunnelJsonFormatException; import java.io.Serializable; import java.math.BigDecimal; +import java.math.BigInteger; import java.time.LocalDate; import java.time.LocalDateTime; import java.time.LocalTime; @@ -112,49 +113,49 @@ public class RowToJsonConverters implements Serializable { return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((byte) value); + return createNumericNode(mapper, value, sqlType); } }; case SMALLINT: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((short) value); + return createNumericNode(mapper, value, sqlType); } }; case INT: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((int) value); + return createNumericNode(mapper, value, sqlType); } }; case BIGINT: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((long) value); + return createNumericNode(mapper, value, sqlType); } }; case FLOAT: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((float) value); + return createNumericNode(mapper, value, sqlType); } }; case DOUBLE: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((double) value); + return createNumericNode(mapper, value, sqlType); } }; case DECIMAL: return new RowToJsonConverter() { @Override public JsonNode convert(ObjectMapper mapper, JsonNode reuse, Object value) { - return mapper.getNodeFactory().numberNode((BigDecimal) value); + return createNumericNode(mapper, value, sqlType); } }; case BYTES: @@ -272,6 +273,59 @@ public class RowToJsonConverters implements Serializable { }; } + /** + * Serializes a value declared as a numeric type in the catalog schema. + * + * <p>Multi-table jobs may match tables whose physical field types differ from the declared + * type, so the runtime representation takes precedence over the declared type instead of being + * blindly cast to it. Only numeric values and character sequences are accepted here: anything + * else still fails fast, so a genuine schema/runtime mismatch is not silently serialized as the + * {@code toString()} of an arbitrary object. + */ + private JsonNode createNumericNode(ObjectMapper mapper, Object value, SqlType declaredType) { + if (value instanceof Byte) { + return mapper.getNodeFactory().numberNode((Byte) value); + } + if (value instanceof Short) { + return mapper.getNodeFactory().numberNode((Short) value); + } + if (value instanceof Integer) { + return mapper.getNodeFactory().numberNode((Integer) value); + } + if (value instanceof Long) { + return mapper.getNodeFactory().numberNode((Long) value); + } + if (value instanceof Float) { + return SqlType.DECIMAL.equals(declaredType) + ? mapper.getNodeFactory().numberNode(BigDecimal.valueOf((Float) value)) + : mapper.getNodeFactory().numberNode((Float) value); + } + if (value instanceof Double) { + return SqlType.DECIMAL.equals(declaredType) + ? mapper.getNodeFactory().numberNode(BigDecimal.valueOf((Double) value)) + : mapper.getNodeFactory().numberNode((Double) value); + } + if (value instanceof BigInteger) { + return mapper.getNodeFactory().numberNode((BigInteger) value); + } + if (value instanceof BigDecimal) { + return mapper.getNodeFactory().numberNode((BigDecimal) value); + } + if (value instanceof CharSequence) { + String text = value.toString(); + try { + return mapper.getNodeFactory().numberNode(new BigDecimal(text)); + } catch (NumberFormatException e) { + return mapper.getNodeFactory().textNode(text); + } + } + throw new SeaTunnelJsonFormatException( + CommonErrorCodeDeprecated.UNSUPPORTED_DATA_TYPE, + String.format( + "Cannot serialize value of type '%s' into the field declared as '%s'", + value.getClass().getName(), declaredType)); + } + private RowToJsonConverter createArrayConverter(ArrayType arrayType) { final RowToJsonConverter elementConverter = createConverter(arrayType.getElementType()); return new RowToJsonConverter() { diff --git a/seatunnel-formats/seatunnel-format-json/src/test/java/org/apache/seatunnel/format/json/JsonRowDataSerDeSchemaTest.java b/seatunnel-formats/seatunnel-format-json/src/test/java/org/apache/seatunnel/format/json/JsonRowDataSerDeSchemaTest.java index 80224658fa..e91070114c 100644 --- a/seatunnel-formats/seatunnel-format-json/src/test/java/org/apache/seatunnel/format/json/JsonRowDataSerDeSchemaTest.java +++ b/seatunnel-formats/seatunnel-format-json/src/test/java/org/apache/seatunnel/format/json/JsonRowDataSerDeSchemaTest.java @@ -294,6 +294,54 @@ public class JsonRowDataSerDeSchemaTest { } } + @Test + public void testSerializeHeterogeneousNumericFields() { + SeaTunnelRowType schema = + new SeaTunnelRowType( + new String[] {"value", "amount"}, + new SeaTunnelDataType[] {INT_TYPE, LONG_TYPE}); + SeaTunnelRow row = new SeaTunnelRow(new Object[] {"text value", new BigDecimal("123.45")}); + + assertEquals( + "{\"value\":\"text value\",\"amount\":123.45}", + new String( + new JsonSerializationSchema(schema).serialize(row), + StandardCharsets.UTF_8)); + } + + @Test + public void testSerializeCrossNumericRuntimeTypes() { + SeaTunnelRowType schema = + new SeaTunnelRowType( + new String[] {"c_int", "c_bigint", "c_float", "c_decimal", "c_str"}, + new SeaTunnelDataType[] { + INT_TYPE, LONG_TYPE, FLOAT_TYPE, new DecimalType(10, 2), INT_TYPE + }); + SeaTunnelRow row = + new SeaTunnelRow( + new Object[] {10L, Integer.valueOf(20), Double.valueOf(1.5D), 2.5D, "123"}); + + assertEquals( + "{\"c_int\":10,\"c_bigint\":20,\"c_float\":1.5,\"c_decimal\":2.5,\"c_str\":123}", + new String( + new JsonSerializationSchema(schema).serialize(row), + StandardCharsets.UTF_8)); + } + + @Test + public void testSerializeNonNumericObjectUnderNumericFieldFails() { + SeaTunnelRowType schema = + new SeaTunnelRowType(new String[] {"c_int"}, new SeaTunnelDataType[] {INT_TYPE}); + SeaTunnelRow row = new SeaTunnelRow(new Object[] {new byte[] {1, 2}}); + + SeaTunnelRuntimeException exception = + Assertions.assertThrows( + SeaTunnelRuntimeException.class, + () -> new JsonSerializationSchema(schema).serialize(row)); + Assertions.assertTrue(exception.getCause() instanceof SeaTunnelJsonFormatException); + Assertions.assertTrue(exception.getCause().getMessage().contains("[B")); + } + @Test public void testSerDeMultiRowsWithNullValues() throws Exception { String[] jsons =
