This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 455d300283 [Fix][Format] Handle heterogeneous numeric JSON fields
(#11415)
455d300283 is described below
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 =