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 =

Reply via email to