This is an automated email from the ASF dual-hosted git repository.
davidzollo 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 9f6b790ff6 [Fix][Connector-V2] Accept any JSON root in Milvus sink
JSON fields (#11526)
9f6b790ff6 is described below
commit 9f6b790ff691832bcc2b1f64f25637693e5a123c
Author: Jast <[email protected]>
AuthorDate: Wed Aug 19 18:45:28 2026 +0800
[Fix][Connector-V2] Accept any JSON root in Milvus sink JSON fields (#11526)
---
.../milvus/utils/sink/MilvusSinkConverter.java | 7 ++-
.../milvus/utils/sink/MilvusSinkConverterTest.java | 56 ++++++++++++++++++++++
2 files changed, 62 insertions(+), 1 deletion(-)
diff --git
a/seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverter.java
b/seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverter.java
index e9f5d26cfc..7d6f03c6d4 100644
---
a/seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverter.java
+++
b/seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverter.java
@@ -73,7 +73,12 @@ public class MilvusSinkConverter {
case STRING:
case DATE:
if (isJson) {
- return gson.fromJson(value.toString(), JsonObject.class);
+ // A Milvus JSON field may hold any JSON root: object,
array or
+ // primitive. Forcing JsonObject fails with "Expected a
+ // com.google.gson.JsonObject but was
com.google.gson.JsonPrimitive"
+ // for non-object values (issue #9677). Object roots still
parse
+ // to JsonObject, so existing behavior is preserved.
+ return JsonParser.parseString(value.toString());
}
return value.toString();
case FLOAT_VECTOR:
diff --git
a/seatunnel-connectors-v2/connector-milvus/src/test/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverterTest.java
b/seatunnel-connectors-v2/connector-milvus/src/test/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverterTest.java
index 95effd3618..7a8dacf411 100644
---
a/seatunnel-connectors-v2/connector-milvus/src/test/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverterTest.java
+++
b/seatunnel-connectors-v2/connector-milvus/src/test/java/org/apache/seatunnel/connectors/seatunnel/milvus/utils/sink/MilvusSinkConverterTest.java
@@ -33,7 +33,10 @@ import
org.apache.seatunnel.connectors.seatunnel.milvus.exception.MilvusConnecto
import org.junit.jupiter.api.Test;
+import com.google.gson.JsonArray;
import com.google.gson.JsonNull;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
import io.milvus.grpc.DataType;
import io.milvus.param.collection.FieldType;
@@ -218,6 +221,59 @@ public class MilvusSinkConverterTest {
assertEquals("MILVUS-10", exception.getSeaTunnelErrorCode().getCode());
}
+ @Test
+ void convertsJsonFieldWithObjectRootToJsonObject() {
+ Object converted =
+ new MilvusSinkConverter()
+ .convertBySeaTunnelType(
+ BasicType.STRING_TYPE, true,
"{\"a\":1,\"b\":\"x\"}");
+
+ assertTrue(converted instanceof JsonObject);
+ JsonObject object = (JsonObject) converted;
+ assertEquals(1, object.get("a").getAsInt());
+ assertEquals("x", object.get("b").getAsString());
+ }
+
+ @Test
+ void convertsJsonFieldWithPrimitiveRootWithoutFailing() {
+ // Issue #9677: a Milvus JSON field may hold any JSON root; forcing
+ // JsonObject throws "Expected a com.google.gson.JsonObject but was
+ // com.google.gson.JsonPrimitive" for non-object values.
+ MilvusSinkConverter converter = new MilvusSinkConverter();
+
+ Object stringRoot =
+ converter.convertBySeaTunnelType(BasicType.STRING_TYPE, true,
"\"abc\"");
+ assertTrue(stringRoot instanceof JsonPrimitive);
+ assertEquals("abc", ((JsonPrimitive) stringRoot).getAsString());
+
+ Object numberRoot =
converter.convertBySeaTunnelType(BasicType.STRING_TYPE, true, "123");
+ assertTrue(numberRoot instanceof JsonPrimitive);
+ assertEquals(123, ((JsonPrimitive) numberRoot).getAsInt());
+
+ Object boolRoot =
converter.convertBySeaTunnelType(BasicType.STRING_TYPE, true, "true");
+ assertTrue(boolRoot instanceof JsonPrimitive);
+ assertTrue(((JsonPrimitive) boolRoot).isBoolean());
+ }
+
+ @Test
+ void convertsJsonFieldWithArrayRootWithoutFailing() {
+ Object converted =
+ new MilvusSinkConverter()
+ .convertBySeaTunnelType(BasicType.STRING_TYPE, true,
"[1,2,3]");
+
+ assertTrue(converted instanceof JsonArray);
+ assertEquals(3, ((JsonArray) converted).size());
+ }
+
+ @Test
+ void keepsNonJsonStringAsIs() {
+ Object converted =
+ new MilvusSinkConverter()
+ .convertBySeaTunnelType(BasicType.STRING_TYPE, false,
"plain");
+
+ assertEquals("plain", converted);
+ }
+
private CatalogTable catalogTable() {
return catalogTable(true);
}