This is an automated email from the ASF dual-hosted git repository.
luchunliang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 0df61b55c4 [INLONG-12177][SDK] JSON format parsing supports returning
JsonObject and JsonArray types (#12178)
0df61b55c4 is described below
commit 0df61b55c4c920e1f385d79abf527419148d36f0
Author: ChunLiang Lu <[email protected]>
AuthorDate: Wed Jul 22 12:42:18 2026 +0800
[INLONG-12177][SDK] JSON format parsing supports returning JsonObject and
JsonArray types (#12178)
* [INLONG-12177][SDK] JSON format parsing now supports returning JsonObject
and JsonArray types
* fix format
---
.../sdk/transform/decode/JsonSourceData.java | 38 ++++++---
.../processor/TestJson2RowDataProcessor.java | 70 +++++++++++++++
.../process/processor/TestKv2RowDataProcessor.java | 99 ++++++++++++++++++++++
3 files changed, 197 insertions(+), 10 deletions(-)
diff --git
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/JsonSourceData.java
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/JsonSourceData.java
index dba0608e4b..43050eeef3 100644
---
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/JsonSourceData.java
+++
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/JsonSourceData.java
@@ -22,6 +22,7 @@ import org.apache.inlong.sdk.transform.process.Context;
import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
import org.apache.commons.lang3.StringUtils;
import java.util.ArrayList;
@@ -73,7 +74,7 @@ public class JsonSourceData extends AbstractSourceData {
* @return
*/
@Override
- public String getField(int rowNum, String fieldName) {
+ public Object getField(int rowNum, String fieldName) {
try {
if (isContextField(fieldName)) {
return getContextField(fieldName);
@@ -86,7 +87,7 @@ public class JsonSourceData extends AbstractSourceData {
}
// parse
if (childNodes.size() == 0) {
- return "";
+ return null;
}
// first node
JsonNode firstNode = childNodes.get(0);
@@ -97,27 +98,27 @@ public class JsonSourceData extends AbstractSourceData {
if (rowNum < childRoot.size()) {
current = childRoot.get(rowNum);
} else {
- return "";
+ return null;
}
} else {
// error data
- return "";
+ return null;
}
if (current == null) {
// error data
- return "";
+ return null;
}
// parse other node
for (int i = 1; i < childNodes.size(); i++) {
JsonNode node = childNodes.get(i);
if (!current.isJsonObject()) {
// error data
- return "";
+ return null;
}
JsonElement newElement =
current.getAsJsonObject().get(node.getName());
if (newElement == null) {
// error data
- return "";
+ return null;
}
// node is not array
if (!node.isArray()) {
@@ -128,12 +129,29 @@ public class JsonSourceData extends AbstractSourceData {
current = getElementFromArray(node, newElement);
if (current == null) {
// error data
- return "";
+ return null;
+ }
+ }
+ if (current.isJsonPrimitive()) {
+ JsonPrimitive jsonPrim = (JsonPrimitive) current;
+ if (jsonPrim.isString()) {
+ return jsonPrim.getAsString();
+ } else if (jsonPrim.isBoolean()) {
+ return jsonPrim.getAsBoolean();
+ } else if (jsonPrim.isNumber()) {
+ return jsonPrim.getAsNumber();
}
+ return jsonPrim.toString();
}
- return current.getAsString();
+ if (current.isJsonNull()) {
+ return null;
+ }
+ if (current.isJsonArray() || current.isJsonObject()) {
+ return current;
+ }
+ return current;
} catch (Exception e) {
- return "";
+ return null;
}
}
diff --git
a/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestJson2RowDataProcessor.java
b/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestJson2RowDataProcessor.java
new file mode 100644
index 0000000000..e3eb86edb2
--- /dev/null
+++
b/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestJson2RowDataProcessor.java
@@ -0,0 +1,70 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.transform.process.processor;
+
+import org.apache.inlong.sdk.transform.decode.SourceDecoderFactory;
+import org.apache.inlong.sdk.transform.encode.SinkEncoderFactory;
+import org.apache.inlong.sdk.transform.pojo.FieldInfo;
+import org.apache.inlong.sdk.transform.pojo.JsonSourceInfo;
+import org.apache.inlong.sdk.transform.pojo.RowDataSinkInfo;
+import org.apache.inlong.sdk.transform.pojo.TransformConfig;
+import org.apache.inlong.sdk.transform.process.TransformProcessor;
+
+import org.apache.flink.table.data.RowData;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.HashMap;
+import java.util.List;
+
+public class TestJson2RowDataProcessor extends AbstractProcessorTestBase {
+
+ @Test
+ public void testJson2RowData() throws Exception {
+ List<FieldInfo> sinkFields = this.getTestFieldList("audit_data_time",
"session_begin_time", "session_id",
+ "business", "product_id", "channel",
+ "agent_id", "archive_p1", "archive_p2",
+ "archive_p3", "archive_p4", "array_field");
+ // sql
+ String transformSql = "select '' as audit_data_time,"
+ + "$root.session_begin_time as session_begin_time,"
+ + "$root.session_id as session_id,"
+ + "$root.business as business,"
+ + "$root.product_id as product_id,"
+ + "$root.channel as channel,"
+ + "$root.agent_id as agent_id,"
+ + "$root.archive_p1 as archive_p1,"
+ + "$root.archive_p2 as archive_p2,"
+ + "$root.archive_p3 as archive_p3,"
+ + "$root.archive_p4 as archive_p4,"
+ + "$root.array_field as array_field from source";
+ // case1
+ TransformProcessor<String, RowData> processor =
TransformProcessor.create(
+ new TransformConfig(transformSql),
+ SourceDecoderFactory.createJsonDecoder(new
JsonSourceInfo("UTF-8", null)),
+ SinkEncoderFactory.createRowEncoder(new
RowDataSinkInfo("UTF-8", sinkFields)));
+ String strJson =
+
"{\"session_id\":\"1782780884\",\"session_begin_time\":\"2026-06-30
08:54:56\",\"business\":\"pay\","
+ +
"\"product_id\":\"1314\",\"channel\":\"todo\",\"agent_id\":\"095d2\",\"archive_p1\":\"money\","
+ +
"\"archive_p2\":\"product\",\"archive_p3\":\"short”\",\"archive_p4\":\"\",\"array_field\":[{\"isArray\":true}]}";
+ List<RowData> output = processor.transform(strJson, new HashMap<>());
+ Assert.assertEquals(1, output.size());
+ Assert.assertEquals(output.get(0).getString(9).toString(), "short”");
+ Assert.assertEquals(output.get(0).getString(11).toString(),
"[{\"isArray\":true}]");
+ }
+}
diff --git
a/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestKv2RowDataProcessor.java
b/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestKv2RowDataProcessor.java
new file mode 100644
index 0000000000..c15eedd084
--- /dev/null
+++
b/inlong-sdk/transform-sdk/src/test/java/org/apache/inlong/sdk/transform/process/processor/TestKv2RowDataProcessor.java
@@ -0,0 +1,99 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sdk.transform.process.processor;
+
+import org.apache.inlong.common.pojo.sort.dataflow.field.format.LongFormatInfo;
+import
org.apache.inlong.common.pojo.sort.dataflow.field.format.TimestampFormatInfo;
+import org.apache.inlong.sdk.transform.decode.SourceDecoderFactory;
+import org.apache.inlong.sdk.transform.encode.SinkEncoderFactory;
+import org.apache.inlong.sdk.transform.pojo.FieldInfo;
+import org.apache.inlong.sdk.transform.pojo.KvSourceInfo;
+import org.apache.inlong.sdk.transform.pojo.RowDataSinkInfo;
+import org.apache.inlong.sdk.transform.pojo.TransformConfig;
+import org.apache.inlong.sdk.transform.process.TransformProcessor;
+
+import org.apache.flink.table.data.RowData;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.HashMap;
+import java.util.List;
+
+public class TestKv2RowDataProcessor extends AbstractProcessorTestBase {
+
+ @Test
+ public void testKv2RowData() throws Exception {
+ List<FieldInfo> sourceFields =
this.getTestFieldList("audit_data_time", "uin", "uuid",
+ "recotime", "abt", "ct",
+ "cv", "from", "itemid",
+ "songtime", "trace",
+ "itemfeature", "userfeature",
+ "extras", "pos");
+ sourceFields.get(0).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(1).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(2).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(3).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(5).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(6).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(7).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(8).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(9).setFormatInfo(new LongFormatInfo());
+ sourceFields.get(14).setFormatInfo(new LongFormatInfo());
+ final KvSourceInfo kvSource = KvSourceInfo.builder().charset("UTF-8")
+ .entryDelimiter('&')
+ .kvDelimiter('=')
+ .escapeChar('\\')
+ .build();
+ List<FieldInfo> sinkFields = this.getTestFieldList("audit_data_time",
+ "tdbank_imp_date", "uin", "uuid",
+ "recotime", "abt", "ct",
+ "cv", "from", "itemid",
+ "songtime", "trace",
+ "itemfeature", "userfeature",
+ "extras", "pos");
+ sinkFields.get(0).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(1).setFormatInfo(new TimestampFormatInfo());
+ sinkFields.get(2).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(3).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(4).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(6).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(7).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(8).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(9).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(10).setFormatInfo(new LongFormatInfo());
+ sinkFields.get(15).setFormatInfo(new LongFormatInfo());
+ // sql
+ String transformSql =
+ "select audit_data_time as
audit_data_time,from_unix_time($ctx.msgTime/1000) as tdbank_imp_date,"
+ + "uin as uin,uuid as uuid,recotime as recotime,abt as
abt,ct as ct,cv as cv,`from` as `from`,"
+ + "itemid as itemid,songtime as songtime,trace as
trace,itemfeature as itemfeature,"
+ + "userfeature as userfeature,extras as extras,pos as
pos from source";
+ // case1
+ TransformProcessor<String, RowData> processor =
TransformProcessor.create(
+ new TransformConfig(transformSql),
+ SourceDecoderFactory.createKvDecoder(kvSource),
+ SinkEncoderFactory.createRowEncoder(new
RowDataSinkInfo("UTF-8", sinkFields)));
+ String strCsv =
+
"pos=19&itemfeature=itemfeature1&uin=1370000000&userfeature=&itemid=620000000&uuid=5800000000"
+ +
"&recotime=1780000000&extras=recall&ct=1&abt=abt1&trace=trace1&cv=200600&from=2&songtime=209";
+ HashMap<String, Object> extParams = new HashMap<>();
+ extParams.put("msgTime", 1784030508000L);
+ List<RowData> output = processor.transform(strCsv, extParams);
+ Assert.assertEquals(1, output.size());
+ }
+}