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 fe603b93f1 [INLONG-12183][SDK] Transform SDK supports $childIndex to 
map the row data subscript. For PB/JSON field mapping, the $root prefix can be 
omitted (#12184)
fe603b93f1 is described below

commit fe603b93f120ee4e0f2ad5ed800d3c48741a5ba5
Author: ChunLiang Lu <[email protected]>
AuthorDate: Tue Aug 18 15:48:04 2026 +0800

    [INLONG-12183][SDK] Transform SDK supports $childIndex to map the row data 
subscript. For PB/JSON field mapping, the $root prefix can be omitted (#12184)
---
 .../sdk/transform/decode/JsonSourceData.java       |  10 +-
 .../inlong/sdk/transform/decode/PbSourceData.java  |   5 +
 .../processor/TestJson2RowDataProcessor.java       | 204 +++++++++++++++++++++
 3 files changed, 217 insertions(+), 2 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 fd229df605..b3dcaeda66 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
@@ -38,6 +38,8 @@ public class JsonSourceData extends AbstractSourceData {
 
     public static final String CHILD_KEY = "$child";
 
+    public static final String CHILD_INDEX_KEY = "$childIndex";
+
     private JsonObject root;
 
     private JsonArray childRoot;
@@ -79,6 +81,9 @@ public class JsonSourceData extends AbstractSourceData {
             if (isContextField(fieldName)) {
                 return getContextField(fieldName);
             }
+            if (StringUtils.equals(CHILD_INDEX_KEY, fieldName)) {
+                return rowNum;
+            }
             JsonElement current = getFieldByElement(rowNum, fieldName);
             if (current == null) {
                 return current;
@@ -126,8 +131,9 @@ public class JsonSourceData extends AbstractSourceData {
                     return null;
                 }
             } else {
-                // error data
-                return null;
+                // default root node
+                current = root;
+                childNodes.add(0, new JsonNode(ROOT_KEY));
             }
             if (current == null) {
                 // error data
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/PbSourceData.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/PbSourceData.java
index 7f6fbcdf85..e35712eb23 100644
--- 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/PbSourceData.java
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/decode/PbSourceData.java
@@ -59,6 +59,8 @@ public class PbSourceData extends AbstractSourceData {
 
     public static final String CHILD_KEY = "$child.";
 
+    public static final String CHILD_INDEX_KEY = "$childIndex";
+
     private Descriptors.Descriptor rootDesc;
 
     private Descriptors.Descriptor childDesc;
@@ -130,6 +132,9 @@ public class PbSourceData extends AbstractSourceData {
             if (isContextField(fieldName)) {
                 return getContextField(fieldName);
             }
+            if (StringUtils.equals(CHILD_INDEX_KEY, fieldName)) {
+                return rowNum;
+            }
             Object fieldValue = findFieldNode(rowNum, fieldName);
             List<PbNode> childNodes = this.columnNodeMap.get(fieldName);
             if (childNodes == null || childNodes.size() == 0) {
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
index f75aacad0b..4897f75014 100644
--- 
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
@@ -986,4 +986,208 @@ public class TestJson2RowDataProcessor extends 
AbstractProcessorTestBase {
         Assert.assertEquals("v_c", personRow.getString(2).toString());
         Assert.assertEquals("v_b", personRow.getString(3).toString());
     }
+
+    // ========== JsonSourceData: $childIndex and omitted $root prefix 
==========
+
+    /**
+     * Verify that `$childIndex` maps to the current row index (0-based) when 
the JSON source
+     * is configured with a child-array root.
+     */
+    @Test
+    public void testJsonChildIndexMapping() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("rowIdx", "sid", 
"msg");
+        // childRoot points to msgs -> multi-row output
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", "msgs");
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select $childIndex as rowIdx, $root.sid as sid, $child.msg as 
msg from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"sid\":\"s1\",\"msgs\":["
+                + "{\"msg\":\"m0\"},"
+                + "{\"msg\":\"m1\"},"
+                + "{\"msg\":\"m2\"}"
+                + "]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(3, output.size());
+
+        // row 0
+        Assert.assertEquals("0", output.get(0).getString(0).toString());
+        Assert.assertEquals("s1", output.get(0).getString(1).toString());
+        Assert.assertEquals("m0", output.get(0).getString(2).toString());
+        // row 1
+        Assert.assertEquals("1", output.get(1).getString(0).toString());
+        Assert.assertEquals("s1", output.get(1).getString(1).toString());
+        Assert.assertEquals("m1", output.get(1).getString(2).toString());
+        // row 2
+        Assert.assertEquals("2", output.get(2).getString(0).toString());
+        Assert.assertEquals("s1", output.get(2).getString(1).toString());
+        Assert.assertEquals("m2", output.get(2).getString(2).toString());
+    }
+
+    /**
+     * Verify that `$childIndex` also works when the source is a single-row 
JSON (no child array),
+     * in which case the index should always be 0.
+     */
+    @Test
+    public void testJsonChildIndexSingleRow() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("rowIdx", "name");
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", null);
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql = "select $childIndex as rowIdx, $root.name as 
name from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"name\":\"Jane\"}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        Assert.assertEquals("0", output.get(0).getString(0).toString());
+        Assert.assertEquals("Jane", output.get(0).getString(1).toString());
+    }
+
+    /**
+     * Verify that the `$root` prefix can be omitted for JSON field mapping.
+     * `session_id` should behave identically to `$root.session_id`.
+     */
+    @Test
+    public void testJsonOmitRootPrefix() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList(
+                "session_id", "business", "product_id", "channel");
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", null);
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // field names without $root prefix
+        String transformSql = "select session_id as session_id,"
+                + "business as business,"
+                + "product_id as product_id,"
+                + "channel as channel from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"session_id\":\"1782780884\","
+                + "\"business\":\"pay\","
+                + "\"product_id\":\"1314\","
+                + "\"channel\":\"todo\"}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        Assert.assertEquals("1782780884", 
output.get(0).getString(0).toString());
+        Assert.assertEquals("pay", output.get(0).getString(1).toString());
+        Assert.assertEquals("1314", output.get(0).getString(2).toString());
+        Assert.assertEquals("todo", output.get(0).getString(3).toString());
+    }
+
+    /**
+     * Verify that omitting `$root` supports nested field paths (e.g. 
`person.name`),
+     * behaving identically to `$root.person.name`.
+     */
+    @Test
+    public void testJsonOmitRootPrefixWithNestedPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("name", "city");
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", null);
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // nested paths without $root prefix
+        String transformSql = "select person.name as name,"
+                + "person.address.city as city from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\","
+                + "\"address\":{\"city\":\"NYC\",\"zip\":\"10001\"}}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        Assert.assertEquals("Jane", output.get(0).getString(0).toString());
+        Assert.assertEquals("NYC", output.get(0).getString(1).toString());
+    }
+
+    /**
+     * Verify that omitting `$root` yields the same result as explicitly using 
`$root`,
+     * and that both styles can be mixed in the same SQL.
+     */
+    @Test
+    public void testJsonOmitRootPrefixMixedWithExplicitRoot() throws Exception 
{
+        List<FieldInfo> sinkFields = this.getTestFieldList("a1", "a2", "b1", 
"b2");
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", null);
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // mix explicit $root and omitted-prefix styles
+        String transformSql = "select $root.name as a1,"
+                + "name as a2,"
+                + "$root.person.name as b1,"
+                + "person.name as b2 from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"name\":\"John\",\"person\":{\"name\":\"Jane\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // a1 and a2 must be identical; b1 and b2 must be identical
+        Assert.assertEquals("John", output.get(0).getString(0).toString());
+        Assert.assertEquals("John", output.get(0).getString(1).toString());
+        Assert.assertEquals("Jane", output.get(0).getString(2).toString());
+        Assert.assertEquals("Jane", output.get(0).getString(3).toString());
+    }
+
+    /**
+     * Combined scenario: within a child-array row source, use both 
`$childIndex`
+     * and omitted-`$root` fields side-by-side.
+     */
+    @Test
+    public void testJsonChildIndexWithOmittedRootPrefix() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("idx", "sid", 
"msg");
+        JsonSourceInfo jsonSource = new JsonSourceInfo("UTF-8", "msgs");
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // `sid` is omitted-$root (root-level field); `$child.msg` is the 
child element field
+        String transformSql =
+                "select $childIndex as idx, sid as sid, $child.msg as msg from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(jsonSource),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"sid\":\"s1\",\"msgs\":["
+                + "{\"msg\":\"m0\"},"
+                + "{\"msg\":\"m1\"}"
+                + "]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(2, output.size());
+        Assert.assertEquals("0", output.get(0).getString(0).toString());
+        Assert.assertEquals("s1", output.get(0).getString(1).toString());
+        Assert.assertEquals("m0", output.get(0).getString(2).toString());
+        Assert.assertEquals("1", output.get(1).getString(0).toString());
+        Assert.assertEquals("s1", output.get(1).getString(1).toString());
+        Assert.assertEquals("m1", output.get(1).getString(2).toString());
+    }
 }

Reply via email to