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 62ff0a944e [INLONG-12181][SDK] Add JSON struct/array extraction 
functions for transform-sdk: json_to_struct, json_to_array, 
json_extract_struct, json_extract_struct_excluding (#12182)
62ff0a944e is described below

commit 62ff0a944eed170c2836d67a462efb6d28bdbc48
Author: ChunLiang Lu <[email protected]>
AuthorDate: Fri Aug 14 11:49:02 2026 +0800

    [INLONG-12181][SDK] Add JSON struct/array extraction functions for 
transform-sdk: json_to_struct, json_to_array, json_extract_struct, 
json_extract_struct_excluding (#12182)
---
 .../sdk/transform/decode/JsonSourceData.java       |  62 +-
 .../sdk/transform/process/TransformProcessor.java  |   6 +
 .../json/JsonExtractStructExcludingFunction.java   | 313 +++++++
 .../function/json/JsonExtractStructFunction.java   | 300 +++++++
 .../process/function/json/JsonToArrayFunction.java | 178 ++++
 .../function/json/JsonToStructFunction.java        | 180 ++++
 .../sdk/transform/process/parser/ArrayParser.java  |  25 +
 .../sdk/transform/utils/FieldToRowDataUtils.java   | 144 ++--
 .../processor/TestJson2RowDataProcessor.java       | 919 +++++++++++++++++++++
 .../process/processor/TestKv2RowDataProcessor.java |  30 +
 10 files changed, 2089 insertions(+), 68 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 43050eeef3..fd229df605 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
@@ -79,12 +79,37 @@ public class JsonSourceData extends AbstractSourceData {
             if (isContextField(fieldName)) {
                 return getContextField(fieldName);
             }
-            // split field name
-            List<JsonNode> childNodes = new ArrayList<>();
-            String[] nodeStrings = fieldName.split("\\.");
-            for (String nodeString : nodeStrings) {
-                childNodes.add(new JsonNode(nodeString));
+            JsonElement current = getFieldByElement(rowNum, fieldName);
+            if (current == null) {
+                return current;
+            }
+            if (current.isJsonNull()) {
+                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();
+            }
+            if (current.isJsonArray() || current.isJsonObject()) {
+                return current;
             }
+            return current;
+        } catch (Exception e) {
+            return null;
+        }
+    }
+
+    public JsonElement getFieldByElement(int rowNum, String fieldName) {
+        try {
+            // split field name
+            List<JsonNode> childNodes = parseNodeList(fieldName);
             // parse
             if (childNodes.size() == 0) {
                 return null;
@@ -132,29 +157,22 @@ public class JsonSourceData extends AbstractSourceData {
                     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();
-            }
-            if (current.isJsonNull()) {
-                return null;
-            }
-            if (current.isJsonArray() || current.isJsonObject()) {
-                return current;
-            }
             return current;
         } catch (Exception e) {
             return null;
         }
     }
 
+    public static List<JsonNode> parseNodeList(String srcFieldName) {
+        // split field name
+        List<JsonNode> childNodes = new ArrayList<>();
+        String[] nodeStrings = srcFieldName.split("\\.");
+        for (String nodeString : nodeStrings) {
+            childNodes.add(new JsonNode(nodeString));
+        }
+        return childNodes;
+    }
+
     private JsonElement getElementFromArray(JsonNode node, JsonElement 
curElement) {
         if (node.getArrayIndices().isEmpty()) {
             // error data
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/TransformProcessor.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/TransformProcessor.java
index f0e39c04bd..74347b31a8 100644
--- 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/TransformProcessor.java
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/TransformProcessor.java
@@ -38,6 +38,8 @@ import net.sf.jsqlparser.statement.select.Select;
 import net.sf.jsqlparser.statement.select.SelectExpressionItem;
 import net.sf.jsqlparser.statement.select.SelectItem;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
 
 import java.io.StringReader;
 import java.util.ArrayList;
@@ -183,6 +185,10 @@ public class TransformProcessor<I, O> {
                     Object fieldValue = parser.parse(sourceData, i, context);
                     if (fieldValue == null) {
                         sinkData.addField(fieldName, "");
+                    } else if (fieldValue instanceof GenericRowData
+                            || fieldValue instanceof GenericArrayData) {
+                        sinkData.addField(fieldName, fieldValue);
+                        context.put(fieldName, fieldValue);
                     } else {
                         sinkData.addField(fieldName, fieldValue.toString());
                         context.put(fieldName, fieldValue);
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructExcludingFunction.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructExcludingFunction.java
new file mode 100644
index 0000000000..d259feb3e8
--- /dev/null
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructExcludingFunction.java
@@ -0,0 +1,313 @@
+/*
+ * 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.function.json;
+
+import org.apache.inlong.sdk.transform.decode.JsonSourceData;
+import org.apache.inlong.sdk.transform.decode.SourceData;
+import org.apache.inlong.sdk.transform.process.Context;
+import org.apache.inlong.sdk.transform.process.function.FunctionConstant;
+import org.apache.inlong.sdk.transform.process.function.TransformFunction;
+import org.apache.inlong.sdk.transform.process.operator.OperatorTools;
+import org.apache.inlong.sdk.transform.process.parser.ColumnParser;
+import org.apache.inlong.sdk.transform.process.parser.ValueParser;
+
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
+import net.sf.jsqlparser.expression.Expression;
+import net.sf.jsqlparser.expression.Function;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * JsonExtractStructExcludingFunction  ->  json_extract_struct_excluding(path, 
excludeField1, excludeField2, ...)
+ * description:
+ * - Only works on JSON source data; returns NULL if the source is not a 
JsonSourceData.
+ * - Returns NULL if 'path' is missing/invalid, or the path cannot be resolved 
to a
+ *   JSON object or JSON array of objects.
+ * - When 'path' resolves to a JSON object, returns a GenericRowData 
containing all
+ *   fields except the specified excluded fields, in the original order.
+ * - When 'path' resolves to a JSON array of objects, returns a 
GenericArrayData whose
+ *   elements are GenericRowData for each array element (all fields except the 
excluded ones).
+ * - Nested json_extract_struct_excluding is supported: e.g.
+ *   json_extract_struct_excluding(json_extract_struct_excluding($root.person, 
address), phone)
+ */
+@TransformFunction(type = FunctionConstant.JSON_TYPE, names = {
+        "json_extract_struct_excluding"}, parameter = "(path, excludeField1, 
excludeField2, ...)", descriptions = {
+                "- Only works on JSON source data; returns NULL if the source 
is not a JsonSourceData;",
+                "- Returns NULL if 'path' is missing/invalid, or the path 
cannot be resolved "
+                        + "to a JSON object or array;",
+                "- When 'path' resolves to a JSON object, returns a 
GenericRowData containing "
+                        + "all fields except the specified excluded fields;",
+                "- When 'path' resolves to a JSON array of objects, returns a 
GenericArrayData "
+                        + "whose elements are GenericRowData for each array 
element (all fields except excluded).",
+                "- Nested json_extract_struct_excluding is supported."
+        }, examples = {
+                "json_extract_struct_excluding($root.person,address,phone) "
+                        + "= <GenericRowData of person without address and 
phone>"
+        })
+public class JsonExtractStructExcludingFunction implements ValueParser {
+
+    private final ValueParser pathParser;
+    private final List<ValueParser> fieldParsers;
+    private String path;
+    private boolean isNestedStruct = false;
+    private boolean isKeepMessage = false;
+
+    public JsonExtractStructExcludingFunction(Function expr) {
+        List<Expression> expressions = expr.getParameters().getExpressions();
+        this.pathParser = OperatorTools.buildParser(expressions.get(0));
+        if (pathParser instanceof ColumnParser) {
+            this.path = ((ColumnParser) pathParser).getFieldName();
+        } else if (pathParser instanceof JsonExtractStructExcludingFunction) {
+            this.isNestedStruct = true;
+        }
+        this.fieldParsers = new ArrayList<>();
+        for (int i = 1; i < expressions.size(); i++) {
+            
this.fieldParsers.add(OperatorTools.buildParser(expressions.get(i)));
+        }
+    }
+
+    @Override
+    public Object parse(SourceData sourceData, int rowIndex, Context context) {
+        if (!(sourceData instanceof JsonSourceData)) {
+            return null;
+        }
+        JsonSourceData jsonData = (JsonSourceData) sourceData;
+
+        // Handle nested json_extract_struct_excluding as path
+        if (isNestedStruct) {
+            JsonExtractStructExcludingFunction child = 
(JsonExtractStructExcludingFunction) pathParser;
+            child.setKeepMessage(true);
+            Object nestedResult = child.parse(sourceData, rowIndex, context);
+            child.setKeepMessage(false);
+            if (nestedResult instanceof JsonObject) {
+                return buildStruct((JsonObject) nestedResult, rowIndex, 
context);
+            }
+            if (nestedResult instanceof JsonArray) {
+                return buildArrayStruct((JsonArray) nestedResult, rowIndex, 
context);
+            }
+            if (nestedResult instanceof GenericArrayData) {
+                return nestedResult;
+            }
+            return null;
+        }
+
+        if (path == null) {
+            return null;
+        }
+
+        // Get the field value at the specified path
+        Object fieldValue = jsonData.getField(rowIndex, path);
+        if (fieldValue == null) {
+            return null;
+        }
+
+        // Handle JsonObject: build a single GenericRowData (all fields except 
excluded)
+        if (fieldValue instanceof JsonObject) {
+            if (isKeepMessage()) {
+                return fieldValue;
+            }
+            return buildStruct((JsonObject) fieldValue, rowIndex, context);
+        }
+
+        // Handle JsonArray: build GenericArrayData of GenericRowData
+        if (fieldValue instanceof JsonArray) {
+            if (isKeepMessage()) {
+                return fieldValue;
+            }
+            return buildArrayStruct((JsonArray) fieldValue, rowIndex, context);
+        }
+
+        // Primitive or null JsonElement: not a struct
+        if (fieldValue instanceof JsonElement) {
+            return null;
+        }
+
+        // Handle GenericRowData (from nested json_extract_struct_excluding)
+        if (fieldValue instanceof GenericRowData) {
+            return fieldValue;
+        }
+
+        // Handle GenericArrayData (from nested json_extract_struct_excluding)
+        if (fieldValue instanceof GenericArrayData) {
+            return fieldValue;
+        }
+
+        return null;
+    }
+
+    /**
+     * Build a GenericRowData from a JsonObject, excluding the specified 
fields.
+     * All fields except the excluded ones are included in the result, in their
+     * original insertion order within the JsonObject.
+     */
+    private GenericRowData buildStruct(JsonObject jsonObject, int rowIndex, 
Context context) {
+        // Collect excluded field names from the column parsers
+        Set<String> excludedFields = new HashSet<>();
+        for (ValueParser parser : fieldParsers) {
+            if (parser instanceof ColumnParser) {
+                excludedFields.add(((ColumnParser) parser).getFieldName());
+            }
+        }
+
+        // Collect non-excluded entries in order
+        List<Map.Entry<String, JsonElement>> includedEntries = new 
ArrayList<>();
+        for (Map.Entry<String, JsonElement> entry : jsonObject.entrySet()) {
+            if (!excludedFields.contains(entry.getKey())) {
+                includedEntries.add(entry);
+            }
+        }
+
+        GenericRowData result = new GenericRowData(includedEntries.size());
+        int index = 0;
+        for (Map.Entry<String, JsonElement> entry : includedEntries) {
+            Object value = convertJsonElement(entry.getValue());
+            result.setField(index++, value);
+        }
+        return result;
+    }
+
+    /**
+     * Build a GenericArrayData from a JsonArray where each element is expected
+     * to be a JsonObject. Each JsonObject is converted to GenericRowData via 
buildStruct,
+     * with the excluded fields removed.
+     */
+    private GenericArrayData buildArrayStruct(JsonArray jsonArray, int 
rowIndex, Context context) {
+        List<Object> valueResult = new ArrayList<>(jsonArray.size());
+        for (int i = 0; i < jsonArray.size(); i++) {
+            JsonElement element = jsonArray.get(i);
+            if (element.isJsonObject()) {
+                valueResult.add(buildStruct(element.getAsJsonObject(), 
rowIndex, context));
+            } else if (element.isJsonNull()) {
+                valueResult.add(null);
+            } else {
+                // Primitive element in array: convert to its Java type
+                valueResult.add(convertJsonElement(element));
+            }
+        }
+        return new GenericArrayData(valueResult.toArray());
+    }
+
+    /**
+     * Convert a JsonElement to its corresponding Java type or Flink data 
structure.
+     * <ul>
+     *   <li>JsonPrimitive String → {@link String}</li>
+     *   <li>JsonPrimitive Boolean → {@link Boolean}</li>
+     *   <li>JsonPrimitive Number → {@link Number}</li>
+     *   <li>JsonObject → the original {@link JsonObject} (for further 
processing)</li>
+     *   <li>JsonArray → {@link GenericArrayData} with each element converted 
recursively</li>
+     *   <li>JsonNull → {@code null}</li>
+     * </ul>
+     */
+    private Object convertJsonElement(JsonElement element) {
+        if (element == null || element.isJsonNull()) {
+            return null;
+        }
+        if (element.isJsonPrimitive()) {
+            JsonPrimitive jsonPrim = element.getAsJsonPrimitive();
+            if (jsonPrim.isString()) {
+                return jsonPrim.getAsString();
+            }
+            if (jsonPrim.isBoolean()) {
+                return jsonPrim.getAsBoolean();
+            }
+            if (jsonPrim.isNumber()) {
+                return jsonPrim.getAsNumber();
+            }
+            return jsonPrim.getAsString();
+        }
+        if (element.isJsonObject()) {
+            // Return the raw JsonObject; caller can extract sub-fields via 
nested json_extract_struct_excluding
+            return element.getAsJsonObject();
+        }
+        if (element.isJsonArray()) {
+            JsonArray jsonArray = element.getAsJsonArray();
+            List<Object> list = new ArrayList<>(jsonArray.size());
+            for (int i = 0; i < jsonArray.size(); i++) {
+                list.add(convertJsonElement(jsonArray.get(i)));
+            }
+            return new GenericArrayData(list.toArray());
+        }
+        return element.toString();
+    }
+
+    /**
+     * Extract a field value from a JsonObject, supporting nested paths
+     * separated by dots (e.g., "address.street").
+     */
+    private Object getJsonFieldValue(JsonObject jsonObject, String fieldName) {
+        String[] parts = fieldName.split("\\.");
+        JsonElement current = jsonObject;
+        for (int i = 0; i < parts.length; i++) {
+            if (current == null || !current.isJsonObject()) {
+                return null;
+            }
+            current = current.getAsJsonObject().get(parts[i]);
+            if (current == null || current.isJsonNull()) {
+                return null;
+            }
+        }
+        return convertJsonElement(current);
+    }
+
+    /**
+     * Build struct data from a raw JsonObject (used internally for nested 
struct processing).
+     */
+    Object buildFromJsonObject(JsonObject jsonObject, int rowIndex, Context 
context) {
+        if (this.path != null) {
+            // If this function has its own path, resolve it relative to the 
given JsonObject
+            Object value = getJsonFieldValue(jsonObject, this.path);
+            if (value instanceof JsonObject) {
+                return buildStruct((JsonObject) value, rowIndex, context);
+            }
+            if (value instanceof JsonArray) {
+                return buildArrayStruct((JsonArray) value, rowIndex, context);
+            }
+            return null;
+        }
+        return buildStruct(jsonObject, rowIndex, context);
+    }
+
+    /**
+     * Check whether the keep-message flag is set.
+     * When true, parse() returns the raw JsonObject/JsonArray instead of 
GenericRowData.
+     *
+     * @return the isKeepMessage
+     */
+    public boolean isKeepMessage() {
+        return isKeepMessage;
+    }
+
+    /**
+     * Set the keep-message flag.
+     * When true, parse() returns the raw JsonObject/JsonArray instead of 
GenericRowData.
+     *
+     * @param isKeepMessage the isKeepMessage to set
+     */
+    public void setKeepMessage(boolean isKeepMessage) {
+        this.isKeepMessage = isKeepMessage;
+    }
+}
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructFunction.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructFunction.java
new file mode 100644
index 0000000000..8c64edb0d3
--- /dev/null
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonExtractStructFunction.java
@@ -0,0 +1,300 @@
+/*
+ * 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.function.json;
+
+import org.apache.inlong.sdk.transform.decode.JsonSourceData;
+import org.apache.inlong.sdk.transform.decode.SourceData;
+import org.apache.inlong.sdk.transform.process.Context;
+import org.apache.inlong.sdk.transform.process.function.FunctionConstant;
+import org.apache.inlong.sdk.transform.process.function.TransformFunction;
+import org.apache.inlong.sdk.transform.process.operator.OperatorTools;
+import org.apache.inlong.sdk.transform.process.parser.ColumnParser;
+import org.apache.inlong.sdk.transform.process.parser.ValueParser;
+
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
+import net.sf.jsqlparser.expression.Expression;
+import net.sf.jsqlparser.expression.Function;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * JsonExtractStructFunction  ->  json_extract_struct(path, field1, field2, 
...)
+ * description:
+ * - Only works on JSON source data; returns NULL if the source is not a 
JsonSourceData.
+ * - Returns NULL if 'path' is missing/invalid, or the path cannot be resolved 
to a
+ *   JSON object or JSON array of objects.
+ * - When 'path' resolves to a JSON object, returns a GenericRowData 
containing the
+ *   specified fields in order. Each field is extracted from the JSON object 
by name;
+ *   fields that cannot be resolved are set to NULL.
+ * - When 'path' resolves to a JSON array of objects, returns a 
GenericArrayData whose
+ *   elements are GenericRowData for each array element.
+ * - Nested json_extract_struct is supported: e.g.
+ *   json_extract_struct(json_extract_struct($root.person, name, age), name)
+ */
+@TransformFunction(type = FunctionConstant.JSON_TYPE, names = {
+        "json_extract_struct"}, parameter = "(path, field1, field2, ...)", 
descriptions = {
+                "- Only works on JSON source data; returns NULL if the source 
is not a JsonSourceData;",
+                "- Returns NULL if 'path' is missing/invalid, or the path 
cannot be resolved "
+                        + "to a JSON object or array;",
+                "- When 'path' resolves to a JSON object, returns a 
GenericRowData containing "
+                        + "the specified fields in order;",
+                "- When 'path' resolves to a JSON array of objects, returns a 
GenericArrayData "
+                        + "whose elements are GenericRowData for each array 
element.",
+                "- Nested json_extract_struct is supported."
+        }, examples = {
+                "json_extract_struct($root.person,name,age) = <GenericRowData 
with fields [name, age]>"
+        })
+public class JsonExtractStructFunction implements ValueParser {
+
+    private final ValueParser pathParser;
+    private final List<ValueParser> fieldParsers;
+    private String path;
+    private boolean isNestedStruct = false;
+    private boolean isKeepMessage = false;
+
+    public JsonExtractStructFunction(Function expr) {
+        List<Expression> expressions = expr.getParameters().getExpressions();
+        this.pathParser = OperatorTools.buildParser(expressions.get(0));
+        if (pathParser instanceof ColumnParser) {
+            this.path = ((ColumnParser) pathParser).getFieldName();
+        } else if (pathParser instanceof JsonExtractStructFunction) {
+            this.isNestedStruct = true;
+        }
+        this.fieldParsers = new ArrayList<>();
+        for (int i = 1; i < expressions.size(); i++) {
+            
this.fieldParsers.add(OperatorTools.buildParser(expressions.get(i)));
+        }
+    }
+
+    @Override
+    public Object parse(SourceData sourceData, int rowIndex, Context context) {
+        if (!(sourceData instanceof JsonSourceData)) {
+            return null;
+        }
+        JsonSourceData jsonData = (JsonSourceData) sourceData;
+
+        // Handle nested json_extract_struct as path
+        if (isNestedStruct) {
+            JsonExtractStructFunction child = (JsonExtractStructFunction) 
pathParser;
+            child.setKeepMessage(true);
+            Object nestedResult = child.parse(sourceData, rowIndex, context);
+            child.setKeepMessage(false);
+            if (nestedResult instanceof JsonObject) {
+                return buildStruct((JsonObject) nestedResult, rowIndex, 
context);
+            }
+            if (nestedResult instanceof JsonArray) {
+                return buildArrayStruct((JsonArray) nestedResult, rowIndex, 
context);
+            }
+            if (nestedResult instanceof GenericArrayData) {
+                return nestedResult;
+            }
+            return null;
+        }
+
+        if (path == null) {
+            return null;
+        }
+
+        // Get the field value at the specified path
+        Object fieldValue = jsonData.getField(rowIndex, path);
+        if (fieldValue == null) {
+            return null;
+        }
+
+        // Handle JsonObject: build a single GenericRowData
+        if (fieldValue instanceof JsonObject) {
+            if (isKeepMessage()) {
+                return fieldValue;
+            }
+            return buildStruct((JsonObject) fieldValue, rowIndex, context);
+        }
+
+        // Handle JsonArray: build GenericArrayData of GenericRowData
+        if (fieldValue instanceof JsonArray) {
+            if (isKeepMessage()) {
+                return fieldValue;
+            }
+            return buildArrayStruct((JsonArray) fieldValue, rowIndex, context);
+        }
+
+        // Primitive or null JsonElement: not a struct
+        if (fieldValue instanceof JsonElement) {
+            return null;
+        }
+
+        // Handle GenericRowData (from nested json_extract_struct)
+        if (fieldValue instanceof GenericRowData) {
+            return fieldValue;
+        }
+
+        // Handle GenericArrayData (from nested json_extract_struct)
+        if (fieldValue instanceof GenericArrayData) {
+            return fieldValue;
+        }
+
+        return null;
+    }
+
+    /**
+     * Build a GenericRowData from a JsonObject using the declared field 
parsers.
+     * Each field parser extracts a value from the JsonObject by its field 
name.
+     * Supports nested paths (e.g., "address.street") by traversing the JSON 
tree.
+     */
+    private GenericRowData buildStruct(JsonObject jsonObject, int rowIndex, 
Context context) {
+        GenericRowData result = new GenericRowData(fieldParsers.size());
+        int index = 0;
+        for (ValueParser parser : fieldParsers) {
+            if (parser instanceof ColumnParser) {
+                ColumnParser columnParser = (ColumnParser) parser;
+                String fieldName = columnParser.getFieldName();
+                Object value = getJsonFieldValue(jsonObject, fieldName);
+                result.setField(index++, value);
+            } else if (parser instanceof JsonExtractStructFunction) {
+                // Nested json_extract_struct as a field value
+                Object value = ((JsonExtractStructFunction) 
parser).buildFromJsonObject(
+                        jsonObject, rowIndex, context);
+                result.setField(index++, value);
+            } else {
+                result.setField(index++, null);
+            }
+        }
+        return result;
+    }
+
+    /**
+     * Build a GenericArrayData from a JsonArray where each element is expected
+     * to be a JsonObject. Each JsonObject is converted to GenericRowData via 
buildStruct.
+     */
+    private GenericArrayData buildArrayStruct(JsonArray jsonArray, int 
rowIndex, Context context) {
+        List<Object> valueResult = new ArrayList<>(jsonArray.size());
+        for (int i = 0; i < jsonArray.size(); i++) {
+            JsonElement element = jsonArray.get(i);
+            if (element.isJsonObject()) {
+                valueResult.add(buildStruct(element.getAsJsonObject(), 
rowIndex, context));
+            } else if (element.isJsonNull()) {
+                valueResult.add(null);
+            } else {
+                // Primitive element in array: convert to its Java type
+                valueResult.add(convertJsonElement(element));
+            }
+        }
+        return new GenericArrayData(valueResult.toArray());
+    }
+
+    /**
+     * Extract a field value from a JsonObject, supporting nested paths
+     * separated by dots (e.g., "address.street").
+     */
+    private Object getJsonFieldValue(JsonObject jsonObject, String fieldName) {
+        String[] parts = fieldName.split("\\.");
+        JsonElement current = jsonObject;
+        for (int i = 0; i < parts.length; i++) {
+            if (current == null || !current.isJsonObject()) {
+                return null;
+            }
+            current = current.getAsJsonObject().get(parts[i]);
+            if (current == null || current.isJsonNull()) {
+                return null;
+            }
+        }
+        return convertJsonElement(current);
+    }
+
+    /**
+     * Convert a JsonElement to its corresponding Java type or Flink data 
structure.
+     * <ul>
+     *   <li>JsonPrimitive String → {@link String}</li>
+     *   <li>JsonPrimitive Boolean → {@link Boolean}</li>
+     *   <li>JsonPrimitive Number → {@link Number}</li>
+     *   <li>JsonObject → the original {@link JsonObject} (for further 
processing)</li>
+     *   <li>JsonArray → {@link GenericArrayData} with each element converted 
recursively</li>
+     *   <li>JsonNull → {@code null}</li>
+     * </ul>
+     */
+    private Object convertJsonElement(JsonElement element) {
+        if (element == null || element.isJsonNull()) {
+            return null;
+        }
+        if (element.isJsonPrimitive()) {
+            JsonPrimitive jsonPrim = element.getAsJsonPrimitive();
+            if (jsonPrim.isString()) {
+                return jsonPrim.getAsString();
+            }
+            if (jsonPrim.isBoolean()) {
+                return jsonPrim.getAsBoolean();
+            }
+            if (jsonPrim.isNumber()) {
+                return jsonPrim.getAsNumber();
+            }
+            return jsonPrim.getAsString();
+        }
+        if (element.isJsonObject()) {
+            // Return the raw JsonObject; caller can extract sub-fields via 
nested json_extract_struct
+            return element.getAsJsonObject();
+        }
+        if (element.isJsonArray()) {
+            JsonArray jsonArray = element.getAsJsonArray();
+            List<Object> list = new ArrayList<>(jsonArray.size());
+            for (int i = 0; i < jsonArray.size(); i++) {
+                list.add(convertJsonElement(jsonArray.get(i)));
+            }
+            return new GenericArrayData(list.toArray());
+        }
+        return element.toString();
+    }
+
+    /**
+     * Build struct data from a raw JsonObject (used internally for nested 
struct processing).
+     */
+    Object buildFromJsonObject(JsonObject jsonObject, int rowIndex, Context 
context) {
+        if (this.path != null) {
+            // If this function has its own path, resolve it relative to the 
given JsonObject
+            Object value = getJsonFieldValue(jsonObject, this.path);
+            if (value instanceof JsonObject) {
+                return buildStruct((JsonObject) value, rowIndex, context);
+            }
+            if (value instanceof JsonArray) {
+                return buildArrayStruct((JsonArray) value, rowIndex, context);
+            }
+            return null;
+        }
+        return buildStruct(jsonObject, rowIndex, context);
+    }
+
+    /**
+     * Check whether the keep-message flag is set.
+     * When true, parse() returns the raw JsonObject/JsonArray instead of 
GenericRowData.
+     */
+    public boolean isKeepMessage() {
+        return isKeepMessage;
+    }
+
+    /**
+     * Set the keep-message flag.
+     * When true, parse() returns the raw JsonObject/JsonArray instead of 
GenericRowData.
+     */
+    public void setKeepMessage(boolean isKeepMessage) {
+        this.isKeepMessage = isKeepMessage;
+    }
+}
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToArrayFunction.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToArrayFunction.java
new file mode 100644
index 0000000000..59fddec454
--- /dev/null
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToArrayFunction.java
@@ -0,0 +1,178 @@
+/*
+ * 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.function.json;
+
+import org.apache.inlong.sdk.transform.decode.JsonSourceData;
+import org.apache.inlong.sdk.transform.decode.SourceData;
+import org.apache.inlong.sdk.transform.process.Context;
+import org.apache.inlong.sdk.transform.process.function.FunctionConstant;
+import org.apache.inlong.sdk.transform.process.function.TransformFunction;
+import org.apache.inlong.sdk.transform.process.operator.OperatorTools;
+import org.apache.inlong.sdk.transform.process.parser.ColumnParser;
+import org.apache.inlong.sdk.transform.process.parser.ValueParser;
+
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
+import net.sf.jsqlparser.expression.Expression;
+import net.sf.jsqlparser.expression.Function;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.binary.BinaryStringData;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * JsonToArrayFunction  ->  json_to_array(path)
+ * description:
+ * - Only works on JSON source data; returns NULL if the source is not a 
JsonSourceData.
+ * - Returns NULL if 'path' is missing/invalid, or the path cannot be resolved 
to a
+ *   JSON array (JsonArray).
+ * - When 'path' resolves to a JSON array, returns a GenericArrayData 
containing the
+ *   complete array. Each element is fully converted:
+ *   <ul>
+ *     <li>JsonObject → GenericRowData with all fields in original order</li>
+ *     <li>JsonPrimitive String → {@link String}</li>
+ *     <li>JsonPrimitive Boolean → {@link Boolean}</li>
+ *     <li>JsonPrimitive Number → {@link Number}</li>
+ *     <li>Nested JsonArray → {@link GenericArrayData} (recursively 
converted)</li>
+ *     <li>JsonNull → {@code null}</li>
+ *   </ul>
+ * - No field filtering is applied; all fields from each element are included.
+ */
+@TransformFunction(type = FunctionConstant.JSON_TYPE, names = {
+        "json_to_array"}, parameter = "(path)", descriptions = {
+                "- Only works on JSON source data; returns NULL if the source 
is not a JsonSourceData;",
+                "- Returns NULL if 'path' is missing/invalid, or the path 
cannot be resolved "
+                        + "to a JSON array (JsonArray);",
+                "- When 'path' resolves to a JSON array, returns a 
GenericArrayData containing "
+                        + "the complete array with all elements fully 
converted.",
+                "- JsonObject elements are converted to GenericRowData with 
all fields preserved.",
+                "- JsonPrimitive elements are converted to their Java types 
(String, Boolean, Number).",
+                "- Nested JsonArray elements are recursively converted to 
GenericArrayData.",
+                "- JsonNull elements are mapped to null.",
+                "- No field filtering is applied."
+        }, examples = {
+                "json_to_array($root.items) = <GenericArrayData of fully 
converted items>"
+        })
+public class JsonToArrayFunction implements ValueParser {
+
+    private final ValueParser pathParser;
+    private String path;
+
+    public JsonToArrayFunction(Function expr) {
+        List<Expression> expressions = expr.getParameters().getExpressions();
+        this.pathParser = OperatorTools.buildParser(expressions.get(0));
+        if (pathParser instanceof ColumnParser) {
+            this.path = ((ColumnParser) pathParser).getFieldName();
+        }
+    }
+
+    @Override
+    public Object parse(SourceData sourceData, int rowIndex, Context context) {
+        if (!(sourceData instanceof JsonSourceData)) {
+            return null;
+        }
+        JsonSourceData jsonData = (JsonSourceData) sourceData;
+
+        if (path == null) {
+            return null;
+        }
+
+        // Get the field value at the specified path
+        Object fieldValue = jsonData.getField(rowIndex, path);
+        if (fieldValue == null) {
+            return null;
+        }
+
+        // Must be a JsonArray, otherwise return null
+        if (!(fieldValue instanceof JsonArray)) {
+            return null;
+        }
+
+        return buildArray((JsonArray) fieldValue);
+    }
+
+    /**
+     * Build a GenericArrayData from a JsonArray.
+     * Each element is fully converted to its corresponding Java type or Flink 
data structure.
+     */
+    private GenericArrayData buildArray(JsonArray jsonArray) {
+        List<Object> valueResult = new ArrayList<>(jsonArray.size());
+        for (int i = 0; i < jsonArray.size(); i++) {
+            JsonElement element = jsonArray.get(i);
+            valueResult.add(convertJsonElement(element));
+        }
+        return new GenericArrayData(valueResult.toArray());
+    }
+
+    /**
+     * Convert a JsonElement to its corresponding Java type or Flink data 
structure.
+     * <ul>
+     *   <li>JsonPrimitive String → {@link String}</li>
+     *   <li>JsonPrimitive Boolean → {@link Boolean}</li>
+     *   <li>JsonPrimitive Number → {@link Number}</li>
+     *   <li>JsonObject → {@link GenericRowData} with all fields preserved</li>
+     *   <li>JsonArray → {@link GenericArrayData} with each element converted 
recursively</li>
+     *   <li>JsonNull → {@code null}</li>
+     * </ul>
+     */
+    private Object convertJsonElement(JsonElement element) {
+        if (element == null || element.isJsonNull()) {
+            return null;
+        }
+        if (element.isJsonPrimitive()) {
+            JsonPrimitive jsonPrim = element.getAsJsonPrimitive();
+            if (jsonPrim.isString()) {
+                return BinaryStringData.fromString(jsonPrim.getAsString());
+            }
+            if (jsonPrim.isBoolean()) {
+                return jsonPrim.getAsBoolean();
+            }
+            if (jsonPrim.isNumber()) {
+                return jsonPrim.getAsNumber();
+            }
+            return BinaryStringData.fromString(jsonPrim.getAsString());
+        }
+        if (element.isJsonObject()) {
+            return buildRow(element.getAsJsonObject());
+        }
+        if (element.isJsonArray()) {
+            return buildArray(element.getAsJsonArray());
+        }
+        return element.toString();
+    }
+
+    /**
+     * Build a GenericRowData from a JsonObject with all fields preserved,
+     * in the original insertion order within the JsonObject.
+     */
+    private GenericRowData buildRow(JsonObject jsonObject) {
+        int fieldCount = jsonObject.size();
+        GenericRowData result = new GenericRowData(fieldCount);
+        int index = 0;
+        for (Map.Entry<String, JsonElement> entry : jsonObject.entrySet()) {
+            Object value = convertJsonElement(entry.getValue());
+            result.setField(index++, value);
+        }
+        return result;
+    }
+}
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToStructFunction.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToStructFunction.java
new file mode 100644
index 0000000000..9da22121d5
--- /dev/null
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/function/json/JsonToStructFunction.java
@@ -0,0 +1,180 @@
+/*
+ * 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.function.json;
+
+import org.apache.inlong.sdk.transform.decode.JsonSourceData;
+import org.apache.inlong.sdk.transform.decode.SourceData;
+import org.apache.inlong.sdk.transform.process.Context;
+import org.apache.inlong.sdk.transform.process.function.FunctionConstant;
+import org.apache.inlong.sdk.transform.process.function.TransformFunction;
+import org.apache.inlong.sdk.transform.process.operator.OperatorTools;
+import org.apache.inlong.sdk.transform.process.parser.ColumnParser;
+import org.apache.inlong.sdk.transform.process.parser.ValueParser;
+
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
+import com.google.gson.JsonPrimitive;
+import net.sf.jsqlparser.expression.Expression;
+import net.sf.jsqlparser.expression.Function;
+import org.apache.flink.table.data.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.binary.BinaryStringData;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * JsonToStructFunction  ->  json_to_struct(path)
+ * description:
+ * - Only works on JSON source data; returns NULL if the source is not a 
JsonSourceData.
+ * - Returns NULL if 'path' is missing/invalid, or the path cannot be resolved 
to a
+ *   JSON object (JsonObject).
+ * - When 'path' resolves to a JSON object, returns a GenericRowData 
containing the
+ *   complete struct. Each field is fully converted:
+ *   <ul>
+ *     <li>JsonPrimitive String → BinaryStringData</li>
+ *     <li>JsonPrimitive Boolean → {@link Boolean}</li>
+ *     <li>JsonPrimitive Number → {@link Number}</li>
+ *     <li>Nested JsonObject → {@link GenericRowData} (recursively 
converted)</li>
+ *     <li>Nested JsonArray → {@link GenericArrayData} (recursively 
converted)</li>
+ *     <li>JsonNull → {@code null}</li>
+ *   </ul>
+ * - No field filtering is applied; all fields from the JsonObject are included
+ *   in original order.
+ */
+@TransformFunction(type = FunctionConstant.JSON_TYPE, names = {
+        "json_to_struct"}, parameter = "(path)", descriptions = {
+                "- Only works on JSON source data; returns NULL if the source 
is not a JsonSourceData;",
+                "- Returns NULL if 'path' is missing/invalid, or the path 
cannot be resolved "
+                        + "to a JSON object (JsonObject);",
+                "- When 'path' resolves to a JSON object, returns a 
GenericRowData containing "
+                        + "the complete struct with all fields fully 
converted.",
+                "- JsonPrimitive elements are converted to their Java types "
+                        + "(String -> BinaryStringData, Boolean, Number).",
+                "- Nested JsonObject elements are recursively converted to 
GenericRowData.",
+                "- Nested JsonArray elements are recursively converted to 
GenericArrayData.",
+                "- JsonNull elements are mapped to null.",
+                "- No field filtering is applied."
+        }, examples = {
+                "json_to_struct($root.person) = <GenericRowData of fully 
converted person object>"
+        })
+public class JsonToStructFunction implements ValueParser {
+
+    private final ValueParser pathParser;
+    private String path;
+
+    public JsonToStructFunction(Function expr) {
+        List<Expression> expressions = expr.getParameters().getExpressions();
+        this.pathParser = OperatorTools.buildParser(expressions.get(0));
+        if (pathParser instanceof ColumnParser) {
+            this.path = ((ColumnParser) pathParser).getFieldName();
+        }
+    }
+
+    @Override
+    public Object parse(SourceData sourceData, int rowIndex, Context context) {
+        if (!(sourceData instanceof JsonSourceData)) {
+            return null;
+        }
+        JsonSourceData jsonData = (JsonSourceData) sourceData;
+
+        if (path == null) {
+            return null;
+        }
+
+        // Get the field value at the specified path
+        Object fieldValue = jsonData.getField(rowIndex, path);
+        if (fieldValue == null) {
+            return null;
+        }
+
+        // Must be a JsonObject, otherwise return null
+        if (!(fieldValue instanceof JsonObject)) {
+            return null;
+        }
+
+        return buildRow((JsonObject) fieldValue);
+    }
+
+    /**
+     * Build a GenericRowData from a JsonObject with all fields preserved,
+     * in the original insertion order within the JsonObject.
+     */
+    private GenericRowData buildRow(JsonObject jsonObject) {
+        int fieldCount = jsonObject.size();
+        GenericRowData result = new GenericRowData(fieldCount);
+        int index = 0;
+        for (Map.Entry<String, JsonElement> entry : jsonObject.entrySet()) {
+            Object value = convertJsonElement(entry.getValue());
+            result.setField(index++, value);
+        }
+        return result;
+    }
+
+    /**
+     * Build a GenericArrayData from a JsonArray.
+     * Each element is fully converted to its corresponding Java type or Flink 
data structure.
+     */
+    private GenericArrayData buildArray(JsonArray jsonArray) {
+        List<Object> valueResult = new ArrayList<>(jsonArray.size());
+        for (int i = 0; i < jsonArray.size(); i++) {
+            JsonElement element = jsonArray.get(i);
+            valueResult.add(convertJsonElement(element));
+        }
+        return new GenericArrayData(valueResult.toArray());
+    }
+
+    /**
+     * Convert a JsonElement to its corresponding Java type or Flink data 
structure.
+     * <ul>
+     *   <li>JsonPrimitive String → BinaryStringData</li>
+     *   <li>JsonPrimitive Boolean → {@link Boolean}</li>
+     *   <li>JsonPrimitive Number → {@link Number}</li>
+     *   <li>JsonObject → {@link GenericRowData} with all fields preserved</li>
+     *   <li>JsonArray → {@link GenericArrayData} with each element converted 
recursively</li>
+     *   <li>JsonNull → {@code null}</li>
+     * </ul>
+     */
+    private Object convertJsonElement(JsonElement element) {
+        if (element == null || element.isJsonNull()) {
+            return null;
+        }
+        if (element.isJsonPrimitive()) {
+            JsonPrimitive jsonPrim = element.getAsJsonPrimitive();
+            if (jsonPrim.isString()) {
+                return BinaryStringData.fromString(jsonPrim.getAsString());
+            }
+            if (jsonPrim.isBoolean()) {
+                return jsonPrim.getAsBoolean();
+            }
+            if (jsonPrim.isNumber()) {
+                return jsonPrim.getAsNumber();
+            }
+            return BinaryStringData.fromString(jsonPrim.getAsString());
+        }
+        if (element.isJsonObject()) {
+            return buildRow(element.getAsJsonObject());
+        }
+        if (element.isJsonArray()) {
+            return buildArray(element.getAsJsonArray());
+        }
+        return BinaryStringData.fromString(element.toString());
+    }
+}
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/parser/ArrayParser.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/parser/ArrayParser.java
index 9571b0a16a..46c2a2689e 100644
--- 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/parser/ArrayParser.java
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/process/parser/ArrayParser.java
@@ -21,6 +21,9 @@ import org.apache.inlong.sdk.transform.decode.SourceData;
 import org.apache.inlong.sdk.transform.process.Context;
 import org.apache.inlong.sdk.transform.process.operator.OperatorTools;
 
+import com.google.gson.JsonArray;
+import com.google.gson.JsonElement;
+import com.google.gson.JsonPrimitive;
 import net.sf.jsqlparser.expression.ArrayExpression;
 
 import java.util.List;
@@ -51,6 +54,28 @@ public class ArrayParser implements ValueParser {
             Number rightObj = (Number) rightValue;
             return leftObj.get(rightObj.intValue());
         }
+        if (leftValue instanceof JsonArray && rightValue instanceof Number) {
+            JsonArray leftObj = (JsonArray) leftValue;
+            Number rightObj = (Number) rightValue;
+            JsonElement result = leftObj.get(rightObj.intValue());
+            if (result.isJsonNull()) {
+                return null;
+            }
+            if (result.isJsonPrimitive()) {
+                JsonPrimitive jsonPrim = (JsonPrimitive) result;
+                if (jsonPrim.isString()) {
+                    return jsonPrim.getAsString();
+                } else if (jsonPrim.isBoolean()) {
+                    return jsonPrim.getAsBoolean();
+                } else if (jsonPrim.isNumber()) {
+                    return jsonPrim.getAsNumber();
+                }
+                return jsonPrim.toString();
+            }
+            if (result.isJsonArray() || result.isJsonObject()) {
+                return result;
+            }
+        }
         return null;
     }
 }
diff --git 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/utils/FieldToRowDataUtils.java
 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/utils/FieldToRowDataUtils.java
index 01d2263020..c204c8f8d5 100644
--- 
a/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/utils/FieldToRowDataUtils.java
+++ 
b/inlong-sdk/transform-sdk/src/main/java/org/apache/inlong/sdk/transform/utils/FieldToRowDataUtils.java
@@ -30,6 +30,7 @@ import org.apache.flink.table.types.logical.DecimalType;
 import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.flink.table.types.logical.LogicalTypeRoot;
 import org.apache.flink.table.types.logical.MapType;
+import org.apache.flink.table.types.logical.RowType;
 
 import java.io.Serializable;
 import java.math.BigDecimal;
@@ -38,7 +39,6 @@ import java.sql.Time;
 import java.sql.Timestamp;
 import java.text.ParseException;
 import java.text.SimpleDateFormat;
-import java.util.Arrays;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -84,9 +84,7 @@ public class FieldToRowDataUtils {
         converterMap.put(LogicalTypeRoot.TIMESTAMP_WITH_TIME_ZONE, (obj) -> 
parseTimestampWithLocalTimeZone(obj));
         converterMap.put(LogicalTypeRoot.DECIMAL, (obj) -> parseDecimal(obj));
         converterMap.put(LogicalTypeRoot.BINARY, (obj) -> parseBinary(obj));
-        converterMap.put(LogicalTypeRoot.ARRAY, (obj) -> parseArray(obj));
         converterMap.put(LogicalTypeRoot.MAP, (obj) -> parseMap(obj));
-        converterMap.put(LogicalTypeRoot.ROW, (obj) -> parseRow(obj));
     }
 
     private static final ThreadLocal<Map<String, SimpleDateFormat>> 
formatLocal = new ThreadLocal<>();
@@ -125,13 +123,54 @@ public class FieldToRowDataUtils {
         }
         switch (type) {
             case ARRAY:
+                final ArrayType aType = (ArrayType) fieldType;
+                final FieldToRowDataConverter elemConverter =
+                        createFieldRowConverter(aType.getElementType());
                 return obj -> {
-                    final Object[] array = (Object[]) obj;
-                    FieldToRowDataConverter elementConverter = 
createFieldRowConverter(
-                            ((ArrayType) fieldType).getElementType());
-                    Object[] converted = Arrays.stream(array)
-                            .map(elementConverter::convert)
-                            .toArray();
+                    Object[] source;
+                    if (obj instanceof GenericArrayData) {
+                        GenericArrayData gArr = (GenericArrayData) obj;
+                        LogicalType elemType = aType.getElementType();
+                        if (elemType.getTypeRoot() == LogicalTypeRoot.ROW) {
+                            int numFields = ((RowType) 
elemType).getFieldCount();
+                            source = new Object[gArr.size()];
+                            for (int i = 0; i < gArr.size(); i++) {
+                                source[i] = gArr.isNullAt(i)
+                                        ? null
+                                        : gArr.getRow(i, numFields);
+                            }
+                        } else {
+                            // Non-ROW element: try to extract and convert 
each element;
+                            // fall back to pass-through if element types 
don't match
+                            try {
+                                Object[] extracted = new Object[gArr.size()];
+                                for (int i = 0; i < gArr.size(); i++) {
+                                    if (gArr.isNullAt(i)) {
+                                        extracted[i] = null;
+                                    } else {
+                                        extracted[i] = extractElement(gArr, i, 
elemType);
+                                    }
+                                }
+                                source = extracted;
+                            } catch (ClassCastException e) {
+                                // Elements don't match declared type (e.g. 
GenericRowData
+                                // in a VARCHAR array); return as-is for 
backward compatibility
+                                return obj;
+                            }
+                        }
+                    } else if (obj instanceof Object[]) {
+                        source = (Object[]) obj;
+                    } else if (obj instanceof List<?>) {
+                        source = ((List<?>) obj).toArray();
+                    } else {
+                        return null;
+                    }
+                    Object[] converted = new Object[source.length];
+                    for (int i = 0; i < source.length; i++) {
+                        converted[i] = source[i] == null
+                                ? null
+                                : elemConverter.convert(source[i]);
+                    }
                     return new GenericArrayData(converted);
                 };
             case MAP:
@@ -148,6 +187,25 @@ public class FieldToRowDataUtils {
                     return new GenericMapData(internalMap);
                 };
             case ROW:
+                final RowType rowType = (RowType) fieldType;
+                final List<RowType.RowField> rowFields = rowType.getFields();
+                final FieldToRowDataConverter[] fieldConverters =
+                        new FieldToRowDataConverter[rowFields.size()];
+                for (int i = 0; i < rowFields.size(); i++) {
+                    fieldConverters[i] = 
createFieldRowConverter(rowFields.get(i).getType());
+                }
+                return obj -> {
+                    if (!(obj instanceof GenericRowData)) {
+                        return null;
+                    }
+                    GenericRowData input = (GenericRowData) obj;
+                    GenericRowData output = new 
GenericRowData(fieldConverters.length);
+                    for (int i = 0; i < fieldConverters.length; i++) {
+                        output.setField(i,
+                                fieldConverters[i].convert(input.getField(i)));
+                    }
+                    return output;
+                };
             case MULTISET:
             case RAW:
             default:
@@ -155,6 +213,37 @@ public class FieldToRowDataUtils {
         }
     }
 
+    private static Object extractElement(GenericArrayData array, int pos, 
LogicalType type) {
+        switch (type.getTypeRoot()) {
+            case ARRAY:
+                return array.getArray(pos);
+            case ROW:
+                return array.getRow(pos, ((RowType) type).getFieldCount());
+            case VARCHAR:
+            case CHAR:
+                return array.getString(pos);
+            case BOOLEAN:
+                return array.getBoolean(pos);
+            case TINYINT:
+                return array.getByte(pos);
+            case SMALLINT:
+                return array.getShort(pos);
+            case INTEGER:
+                return array.getInt(pos);
+            case BIGINT:
+                return array.getLong(pos);
+            case FLOAT:
+                return array.getFloat(pos);
+            case DOUBLE:
+                return array.getDouble(pos);
+            case BINARY:
+            case VARBINARY:
+                return array.getBinary(pos);
+            default:
+                return array.getBinary(pos);
+        }
+    }
+
     private static Object parseBoolean(Object obj) {
         try {
             if (obj == null) {
@@ -423,26 +512,6 @@ public class FieldToRowDataUtils {
         }
     }
 
-    private static Object parseArray(Object obj) {
-        try {
-            if (obj == null) {
-                return null;
-            }
-            if (obj instanceof GenericArrayData) {
-                return obj;
-            }
-            if (obj instanceof List<?>) {
-                return new GenericArrayData(((List<?>) obj).toArray());
-            }
-            return null;
-        } catch (RuntimeException e) {
-            if (isIgnoreError()) {
-                return null;
-            }
-            throw e;
-        }
-    }
-
     private static Object parseMap(Object obj) {
         try {
             if (obj == null) {
@@ -462,21 +531,4 @@ public class FieldToRowDataUtils {
             throw e;
         }
     }
-
-    private static Object parseRow(Object obj) {
-        try {
-            if (obj == null) {
-                return null;
-            }
-            if (obj instanceof GenericRowData) {
-                return obj;
-            }
-            return null;
-        } catch (RuntimeException e) {
-            if (isIgnoreError()) {
-                return null;
-            }
-            throw e;
-        }
-    }
 }
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 e3eb86edb2..f75aacad0b 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
@@ -17,6 +17,10 @@
 
 package org.apache.inlong.sdk.transform.process.processor;
 
+import 
org.apache.inlong.common.pojo.sort.dataflow.field.format.ArrayFormatInfo;
+import org.apache.inlong.common.pojo.sort.dataflow.field.format.FormatInfo;
+import org.apache.inlong.common.pojo.sort.dataflow.field.format.RowFormatInfo;
+import 
org.apache.inlong.common.pojo.sort.dataflow.field.format.StringFormatInfo;
 import org.apache.inlong.sdk.transform.decode.SourceDecoderFactory;
 import org.apache.inlong.sdk.transform.encode.SinkEncoderFactory;
 import org.apache.inlong.sdk.transform.pojo.FieldInfo;
@@ -25,6 +29,8 @@ 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.GenericArrayData;
+import org.apache.flink.table.data.GenericRowData;
 import org.apache.flink.table.data.RowData;
 import org.junit.Assert;
 import org.junit.Test;
@@ -67,4 +73,917 @@ public class TestJson2RowDataProcessor extends 
AbstractProcessorTestBase {
         Assert.assertEquals(output.get(0).getString(9).toString(), "short”");
         Assert.assertEquals(output.get(0).getString(11).toString(), 
"[{\"isArray\":true}]");
     }
+
+    @Test
+    public void testJsonExtractStructWithObject() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("name", "age");
+        // person struct sink field
+        FieldInfo personStruct = new FieldInfo("personStruct");
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+        sinkFields.add(personStruct);
+        // items struct sink field (array of struct)
+        FieldInfo itemsStruct = new FieldInfo("itemsStruct");
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo itemsStructFormat = new 
ArrayFormatInfo(itemsRowFormat);
+        itemsStruct.setFormatInfo(itemsStructFormat);
+        sinkFields.add(itemsStruct);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql = "select $root.name as name,"
+                + "$root.age as age,"
+                + "json_extract_struct($root.person,name,age) as personStruct,"
+                + "json_extract_struct($root.items,id,value) as itemsStruct"
+                + " from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"name\":\"John\",\"age\":30,"
+                + "\"person\":{\"name\":\"Jane\",\"age\":25},"
+                + 
"\"items\":[{\"id\":\"1\",\"value\":\"item1\"},{\"id\":\"2\",\"value\":\"item2\"}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // verify basic fields
+        Assert.assertEquals("John", output.get(0).getString(0).toString());
+        Assert.assertEquals("30", output.get(0).getString(1).toString());
+
+        // verify personStruct: GenericRowData with name=Jane, age=25
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(2, 2);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+
+        // verify itemsStruct: GenericArrayData of GenericRowData
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(3);
+        Assert.assertEquals(2, itemsArray.size());
+        RowData item0 = itemsArray.getRow(0, 2);
+        Assert.assertEquals("1", item0.getString(0).toString());
+        Assert.assertEquals("item1", item0.getString(1).toString());
+        RowData item1 = itemsArray.getRow(1, 2);
+        Assert.assertEquals("2", item1.getString(0).toString());
+        Assert.assertEquals("item2", item1.getString(1).toString());
+    }
+
+    @Test
+    public void testJsonExtractStructWithMissingField() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age", "missing_field"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), new StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select 
json_extract_struct($root.person,name,age,missing_field) as personStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 3);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+        Assert.assertNull(personRow.getString(2)); // missing field returns 
null
+    }
+
+    @Test
+    public void testJsonExtractStructWithNestedPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("addressStruct");
+        FieldInfo addressStruct = sinkFields.get(0);
+        RowFormatInfo addressStructFormat = new RowFormatInfo(
+                new String[]{"city", "zip"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        addressStruct.setFormatInfo(addressStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct($root.address,city,zip) as 
addressStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"address\":{\"city\":\"NYC\",\"zip\":\"10001\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData addressRow = (GenericRowData) output.get(0).getRow(0, 
2);
+        Assert.assertEquals("NYC", addressRow.getString(0).toString());
+        Assert.assertEquals("10001", addressRow.getString(1).toString());
+    }
+
+    @Test
+    public void testJsonExtractStructWithNonExistentPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct($root.non_existent_path,name,age) 
as personStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // non-existent path should return null struct
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonExtractStructWithEmptyArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("itemsStruct");
+        FieldInfo itemsStruct = sinkFields.get(0);
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo itemsStructFormat = new 
ArrayFormatInfo(itemsRowFormat);
+        itemsStruct.setFormatInfo(itemsStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct($root.items,id,value) as 
itemsStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"items\":[]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(0, itemsArray.size());
+    }
+
+    @Test
+    public void testJsonExtractStructWithPrimitiveArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("tagsStruct");
+        FieldInfo tagsStruct = sinkFields.get(0);
+        ArrayFormatInfo tagsStructFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        tagsStruct.setFormatInfo(tagsStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct($root.data,tag) as tagsStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"data\":[{\"tag\":\"a\"},{\"tag\":\"b\"},{\"tag\":\"c\"}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData tagsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(3, tagsArray.size());
+    }
+
+    // ========== JsonExtractStructExcludingFunction tests ==========
+
+    @Test
+    public void testJsonExtractStructExcludingWithObject() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("name", "age");
+        // personStruct sink field: after excluding email,phone from person 
object,
+        // the remaining fields are name,age in original JSON order
+        FieldInfo personStruct = new FieldInfo("personStruct");
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+        sinkFields.add(personStruct);
+        // itemsStruct sink field: after excluding extra from items array 
elements,
+        // the remaining fields are id,value
+        FieldInfo itemsStruct = new FieldInfo("itemsStruct");
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo itemsStructFormat = new 
ArrayFormatInfo(itemsRowFormat);
+        itemsStruct.setFormatInfo(itemsStructFormat);
+        sinkFields.add(itemsStruct);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql = "select $root.name as name,"
+                + "$root.age as age,"
+                + "json_extract_struct_excluding($root.person,email,phone) as 
personStruct,"
+                + "json_extract_struct_excluding($root.items,extra) as 
itemsStruct"
+                + " from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"name\":\"John\",\"age\":30,"
+                + 
"\"person\":{\"name\":\"Jane\",\"age\":25,\"email\":\"[email protected]\",\"phone\":\"123456\"},"
+                + 
"\"items\":[{\"id\":\"1\",\"value\":\"item1\",\"extra\":\"x1\"},"
+                + "{\"id\":\"2\",\"value\":\"item2\",\"extra\":\"x2\"}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // verify basic fields
+        Assert.assertEquals("John", output.get(0).getString(0).toString());
+        Assert.assertEquals("30", output.get(0).getString(1).toString());
+
+        // verify personStruct: GenericRowData with name=Jane, age=25 
(email,phone excluded)
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(2, 2);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+
+        // verify itemsStruct: GenericArrayData of GenericRowData (extra 
excluded from each)
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(3);
+        Assert.assertEquals(2, itemsArray.size());
+        RowData item0 = itemsArray.getRow(0, 2);
+        Assert.assertEquals("1", item0.getString(0).toString());
+        Assert.assertEquals("item1", item0.getString(1).toString());
+        RowData item1 = itemsArray.getRow(1, 2);
+        Assert.assertEquals("2", item1.getString(0).toString());
+        Assert.assertEquals("item2", item1.getString(1).toString());
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithNoExcludedFields() throws 
Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        // No fields excluded, all fields (name,age,email) should be returned 
in order
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age", "email"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), new StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select 
json_extract_struct_excluding($root.person,non_existent_field) as personStruct 
from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"person\":{\"name\":\"Jane\",\"age\":25,\"email\":\"[email protected]\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // All three fields should be present since non_existent_field doesn't 
match
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 3);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+        Assert.assertEquals("[email protected]", 
personRow.getString(2).toString());
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithAllFieldsExcluded() throws 
Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        // All fields excluded, result is an empty GenericRowData
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{},
+                new FormatInfo[]{});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct_excluding($root.person,name,age) 
as personStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // All fields excluded, should return empty GenericRowData with 0 arity
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 0);
+        Assert.assertEquals(0, personRow.getArity());
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("itemsStruct");
+        FieldInfo itemsStruct = sinkFields.get(0);
+        // exclude 'extra' from each array element, remaining: id,value
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo itemsStructFormat = new 
ArrayFormatInfo(itemsRowFormat);
+        itemsStruct.setFormatInfo(itemsStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct_excluding($root.items,extra) as 
itemsStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"items\":[{\"id\":\"1\",\"value\":\"v1\",\"extra\":\"x1\"},"
+                + 
"{\"id\":\"2\",\"value\":\"v2\",\"extra\":\"x2\"},{\"id\":\"3\",\"value\":\"v3\",\"extra\":\"x3\"}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(3, itemsArray.size());
+        RowData item0 = itemsArray.getRow(0, 2);
+        Assert.assertEquals("1", item0.getString(0).toString());
+        Assert.assertEquals("v1", item0.getString(1).toString());
+        RowData item2 = itemsArray.getRow(2, 2);
+        Assert.assertEquals("3", item2.getString(0).toString());
+        Assert.assertEquals("v3", item2.getString(1).toString());
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithEmptyArray() throws 
Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("itemsStruct");
+        FieldInfo itemsStruct = sinkFields.get(0);
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo itemsStructFormat = new 
ArrayFormatInfo(itemsRowFormat);
+        itemsStruct.setFormatInfo(itemsStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct_excluding($root.items,extra) as 
itemsStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"items\":[]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // Empty array should return empty GenericArrayData
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(0, itemsArray.size());
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithNonExistentPath() throws 
Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select 
json_extract_struct_excluding($root.non_existent_path,address) as personStruct 
from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // Non-existent path should return null struct
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonExtractStructExcludingWithObjectPreservesFieldOrder() 
throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        // Original JSON field order: a,d,c,b -> exclude 'd' -> remaining: 
a,c,b (in original order)
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"a", "c", "b"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), new StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_extract_struct_excluding($root.person,d) as 
personStruct from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        // JSON with fields in order: a,d,c,b
+        String strJson = 
"{\"person\":{\"a\":\"v_a\",\"d\":\"v_d\",\"c\":\"v_c\",\"b\":\"v_b\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // After excluding 'd', fields should remain in order: a, c, b
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 3);
+        Assert.assertEquals("v_a", personRow.getString(0).toString());
+        Assert.assertEquals("v_c", personRow.getString(1).toString());
+        Assert.assertEquals("v_b", personRow.getString(2).toString());
+    }
+
+    // ========== JsonToArrayFunction tests ==========
+
+    @Test
+    public void testJsonToArrayWithObjectArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("itemsArray");
+        FieldInfo itemsArrayField = sinkFields.get(0);
+        // Each element has id,name,active in order
+        RowFormatInfo itemsRowFormat = new RowFormatInfo(
+                new String[]{"id", "name", "active"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), new StringFormatInfo()});
+        ArrayFormatInfo itemsArrayFormat = new ArrayFormatInfo(itemsRowFormat);
+        itemsArrayField.setFormatInfo(itemsArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.items) as itemsArray from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"items\":[{\"id\":\"1\",\"name\":\"item1\",\"active\":true},"
+                + "{\"id\":\"2\",\"name\":\"item2\",\"active\":false}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData itemsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(2, itemsArray.size());
+
+        RowData item0 = itemsArray.getRow(0, 3);
+        Assert.assertEquals("1", item0.getString(0).toString());
+        Assert.assertEquals("item1", item0.getString(1).toString());
+        Assert.assertEquals("true", item0.getString(2).toString());
+
+        RowData item1 = itemsArray.getRow(1, 3);
+        Assert.assertEquals("2", item1.getString(0).toString());
+        Assert.assertEquals("item2", item1.getString(1).toString());
+        Assert.assertEquals("false", item1.getString(2).toString());
+    }
+
+    @Test
+    public void testJsonToArrayWithPrimitiveArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("strArray");
+        FieldInfo strArrayField = sinkFields.get(0);
+        ArrayFormatInfo strArrayFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        strArrayField.setFormatInfo(strArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.tags) as strArray from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"tags\":[\"a\",\"b\",\"c\"]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData tagsArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(3, tagsArray.size());
+        Assert.assertEquals("a", tagsArray.getString(0).toString());
+        Assert.assertEquals("b", tagsArray.getString(1).toString());
+        Assert.assertEquals("c", tagsArray.getString(2).toString());
+    }
+
+    @Test
+    public void testJsonToArrayWithNumberArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("numArray");
+        FieldInfo numArrayField = sinkFields.get(0);
+        ArrayFormatInfo numArrayFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        numArrayField.setFormatInfo(numArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.values) as numArray from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"values\":[1,2,3]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData numArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(3, numArray.size());
+    }
+
+    @Test
+    public void testJsonToArrayWithNonArrayPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultArray");
+        FieldInfo resultArrayField = sinkFields.get(0);
+        ArrayFormatInfo resultArrayFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        resultArrayField.setFormatInfo(resultArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // path resolves to a JSON object, not an array → should return null
+        String transformSql =
+                "select json_to_array($root.person) as resultArray from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // Person is an object, not an array → null
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonToArrayWithNonExistentPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultArray");
+        FieldInfo resultArrayField = sinkFields.get(0);
+        ArrayFormatInfo resultArrayFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        resultArrayField.setFormatInfo(resultArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.non_existent) as resultArray from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"items\":[1,2,3]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // Non-existent path → null
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonToArrayWithEmptyArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultArray");
+        FieldInfo resultArrayField = sinkFields.get(0);
+        ArrayFormatInfo resultArrayFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        resultArrayField.setFormatInfo(resultArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.items) as resultArray from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"items\":[]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData resultArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(0, resultArray.size());
+    }
+
+    @Test
+    public void testJsonToArrayWithNestedArray() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("nestedArray");
+        FieldInfo nestedArrayField = sinkFields.get(0);
+        RowFormatInfo innerRowFormat = new RowFormatInfo(
+                new String[]{"id", "value"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        ArrayFormatInfo innerArrayFormat = new ArrayFormatInfo(innerRowFormat);
+        ArrayFormatInfo nestedArrayFormat = new 
ArrayFormatInfo(innerArrayFormat);
+        nestedArrayField.setFormatInfo(nestedArrayFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_array($root.matrix) as nestedArray from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"matrix\":[[{\"id\":\"1\",\"value\":\"a\"}],[{\"id\":\"2\",\"value\":\"b\"}]]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericArrayData outerArray = (GenericArrayData) 
output.get(0).getArray(0);
+        Assert.assertEquals(2, outerArray.size());
+        GenericArrayData innerArray0 = (GenericArrayData) 
outerArray.getArray(0);
+        Assert.assertEquals(1, innerArray0.size());
+        RowData innerItem0 = innerArray0.getRow(0, 2);
+        Assert.assertEquals("1", innerItem0.getString(0).toString());
+        Assert.assertEquals("a", innerItem0.getString(1).toString());
+    }
+
+    // ========== JsonToStructFunction tests ==========
+
+    @Test
+    public void testJsonToStructWithObject() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age", "email"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), new StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.person) as personStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"person\":{\"name\":\"Jane\",\"age\":25,\"email\":\"[email protected]\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 3);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+        Assert.assertEquals("[email protected]", 
personRow.getString(2).toString());
+    }
+
+    @Test
+    public void testJsonToStructWithNonObjectPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultStruct");
+        FieldInfo resultStruct = sinkFields.get(0);
+        RowFormatInfo resultStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        resultStruct.setFormatInfo(resultStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        // path resolves to a JSON array, not an object -> should return null
+        String transformSql =
+                "select json_to_struct($root.items) as resultStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"items\":[{\"name\":\"a\",\"age\":1}]}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // Items is an array, not an object -> null
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonToStructWithNonExistentPath() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultStruct");
+        FieldInfo resultStruct = sinkFields.get(0);
+        RowFormatInfo resultStructFormat = new RowFormatInfo(
+                new String[]{"name", "age"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        resultStruct.setFormatInfo(resultStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.non_existent) as resultStruct 
from source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+        // Non-existent path -> null
+        Assert.assertTrue(output.get(0).isNullAt(0));
+    }
+
+    @Test
+    public void testJsonToStructWithNestedObject() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        // person has name, age, address (nested object with city, zip)
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo addressRowFormat = new RowFormatInfo(
+                new String[]{"city", "zip"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo()});
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "age", "address"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(), addressRowFormat});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.person) as personStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"person\":{\"name\":\"Jane\",\"age\":25,"
+                + "\"address\":{\"city\":\"NYC\",\"zip\":\"10001\"}}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 3);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+        Assert.assertEquals("25", personRow.getString(1).toString());
+
+        // address is a nested GenericRowData
+        GenericRowData addressRow = (GenericRowData) personRow.getRow(2, 2);
+        Assert.assertEquals("NYC", addressRow.getString(0).toString());
+        Assert.assertEquals("10001", addressRow.getString(1).toString());
+    }
+
+    @Test
+    public void testJsonToStructWithEmptyObject() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("resultStruct");
+        FieldInfo resultStruct = sinkFields.get(0);
+        RowFormatInfo resultStructFormat = new RowFormatInfo(
+                new String[]{},
+                new FormatInfo[]{});
+        resultStruct.setFormatInfo(resultStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.empty) as resultStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = "{\"empty\":{}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData resultRow = (GenericRowData) output.get(0).getRow(0, 0);
+        Assert.assertEquals(0, resultRow.getArity());
+    }
+
+    @Test
+    public void testJsonToStructWithArrayField() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        ArrayFormatInfo tagsFormat = new ArrayFormatInfo(new 
StringFormatInfo());
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"name", "tags"},
+                new FormatInfo[]{new StringFormatInfo(), tagsFormat});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.person) as personStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        String strJson = 
"{\"person\":{\"name\":\"Jane\",\"tags\":[\"a\",\"b\",\"c\"]}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 2);
+        Assert.assertEquals("Jane", personRow.getString(0).toString());
+
+        GenericArrayData tagsArray = (GenericArrayData) personRow.getArray(1);
+        Assert.assertEquals(3, tagsArray.size());
+        Assert.assertEquals("a", tagsArray.getString(0).toString());
+        Assert.assertEquals("b", tagsArray.getString(1).toString());
+        Assert.assertEquals("c", tagsArray.getString(2).toString());
+    }
+
+    @Test
+    public void testJsonToStructPreservesFieldOrder() throws Exception {
+        List<FieldInfo> sinkFields = this.getTestFieldList("personStruct");
+        FieldInfo personStruct = sinkFields.get(0);
+        RowFormatInfo personStructFormat = new RowFormatInfo(
+                new String[]{"a", "d", "c", "b"},
+                new FormatInfo[]{new StringFormatInfo(), new 
StringFormatInfo(),
+                        new StringFormatInfo(), new StringFormatInfo()});
+        personStruct.setFormatInfo(personStructFormat);
+
+        RowDataSinkInfo rowSink = new RowDataSinkInfo("UTF-8", sinkFields);
+
+        String transformSql =
+                "select json_to_struct($root.person) as personStruct from 
source";
+
+        TransformConfig config = new TransformConfig(transformSql);
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                config,
+                SourceDecoderFactory.createJsonDecoder(new 
JsonSourceInfo("UTF-8", null)),
+                SinkEncoderFactory.createRowEncoder(rowSink));
+
+        // JSON with fields in order: a,d,c,b
+        String strJson = 
"{\"person\":{\"a\":\"v_a\",\"d\":\"v_d\",\"c\":\"v_c\",\"b\":\"v_b\"}}";
+
+        List<RowData> output = processor.transform(strJson, new HashMap<>());
+        Assert.assertEquals(1, output.size());
+
+        // Fields should remain in original JSON order: a, d, c, b
+        GenericRowData personRow = (GenericRowData) output.get(0).getRow(0, 4);
+        Assert.assertEquals("v_a", personRow.getString(0).toString());
+        Assert.assertEquals("v_d", personRow.getString(1).toString());
+        Assert.assertEquals("v_c", personRow.getString(2).toString());
+        Assert.assertEquals("v_b", personRow.getString(3).toString());
+    }
 }
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
index c15eedd084..5e7753964f 100644
--- 
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
@@ -18,6 +18,7 @@
 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.StringFormatInfo;
 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;
@@ -96,4 +97,33 @@ public class TestKv2RowDataProcessor extends 
AbstractProcessorTestBase {
         List<RowData> output = processor.transform(strCsv, extParams);
         Assert.assertEquals(1, output.size());
     }
+
+    @Test
+    public void testKv2RowData4Base64() throws Exception {
+        List<FieldInfo> sourceFields = this.getTestFieldList("content");
+        sourceFields.get(0).setFormatInfo(new StringFormatInfo());
+        final KvSourceInfo kvSource = KvSourceInfo.builder().charset("UTF-8")
+                .entryDelimiter('&')
+                .kvDelimiter('=')
+                .escapeChar('\\')
+                .build();
+        List<FieldInfo> sinkFields = this.getTestFieldList("content",
+                "decoded_content");
+        sinkFields.get(0).setFormatInfo(new StringFormatInfo());
+        sinkFields.get(1).setFormatInfo(new StringFormatInfo());
+        // sql
+        String transformSql =
+                "select content as content,from_base64(content) as 
decoded_content from source";
+        // case1
+        TransformProcessor<String, RowData> processor = 
TransformProcessor.create(
+                new TransformConfig(transformSql),
+                SourceDecoderFactory.createKvDecoder(kvSource),
+                SinkEncoderFactory.createRowEncoder(new 
RowDataSinkInfo("UTF-8", sinkFields)));
+        String strCsv =
+                
"content=8J+OicKg5Zeo77yBUvCfjY5uZ1LCoArlvojpq5jlhbTorqTor4bkvaDvvIHmiJHmmK/kvaDnmoTkuJPlsZ7ohb7orq/op4bpopHnpo/liKnlrpjvvIzotoXlpJrnpo/liKnlt7LlpIflpb3wn5GHCgrwn5SU5oqi5YWI5p2D55uK772c5YaF5a655LiK5paw44CB5piO5pif5Yqo5oCB5o+Q5YmN6YCa55+lCvCfjJ/kuJPlsZ7npo/liKnvvZznrb7lkI3nhafjgIHnur/kuIvmtLvliqjkuJPlsZ7pgJrpgZMK8J+SsOi2heWAvOS8mOaDoO+9nOS8muWRmOeJueaDoOOAgeWFjei0ueaKveixquekvAoK8J+OgeingemdouWFiOmAgeS9oOS4gOS7veaWsOS6uuekvApodHRwczovL3Z1cmwucXEuY29tL1Q4NVhVMFNHCgrigJTigJT
 [...]
+        HashMap<String, Object> extParams = new HashMap<>();
+        extParams.put("msgTime", 1784030508000L);
+        List<RowData> output = processor.transform(strCsv, extParams);
+        Assert.assertEquals(1, output.size());
+    }
 }

Reply via email to