This is an automated email from the ASF dual-hosted git repository.

dockerzhang 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 254cb3dcb0 [INLONG-8345][Manager][Sort] Support Flink multi-version 
adaptation directory adjustment and version unified management (#8388)
254cb3dcb0 is described below

commit 254cb3dcb0650ad8aaea1c99b7f80e2b837c1a85
Author: ganfengtan <[email protected]>
AuthorDate: Mon Jul 3 20:54:32 2023 +0800

    [INLONG-8345][Manager][Sort] Support Flink multi-version adaptation 
directory adjustment and version unified management (#8388)
---
 inlong-manager/manager-plugins/pom.xml             |   2 +-
 inlong-sort/sort-dist/pom.xml                      |  21 +--
 inlong-sort/sort-flink/base/pom.xml                |   8 +-
 .../sort/base/dirty/sink/log/LogDirtySink.java     |   7 +-
 .../inlong/sort/base/dirty/utils/FormatUtils.java  |   8 +-
 .../format/DebeziumJsonDynamicSchemaFormat.java    |  25 +---
 .../sort/base/format/JsonDynamicSchemaFormat.java  |  88 +----------
 .../sort/base/format/JsonToRowDataConverters.java  |   3 +-
 .../inlong/sort/base/sink/MultipleSinkOption.java  |  10 +-
 .../inlong/sort/base/dirty/FormatUtilsTest.java    |   9 +-
 .../{format-json => format-json-v1.13}/pom.xml     |   4 +-
 .../inlong/sort/formats/json/MysqlBinLogData.java  |   0
 .../inlong/sort/formats/json/canal/CanalJson.java  |   0
 .../json/canal/CanalJsonDecodingFormat.java        |   0
 .../json/canal/CanalJsonDeserializationSchema.java |   0
 .../canal/CanalJsonEnhancedDecodingFormat.java     |   0
 .../CanalJsonEnhancedDeserializationSchema.java    |   0
 .../canal/CanalJsonEnhancedEncodingFormat.java     |   0
 .../json/canal/CanalJsonEnhancedFormatFactory.java |   0
 .../CanalJsonEnhancedSerializationSchema.java      |   0
 .../json/canal/CanalJsonSerializationSchema.java   |   0
 .../inlong/sort/formats/json/canal/CanalUtils.java |   0
 .../sort/formats/json/debezium/DebeziumJson.java   |   0
 .../json/debezium/DebeziumJsonDecodingFormat.java  |   0
 .../DebeziumJsonDeserializationSchema.java         |   0
 .../sort/formats/json/debezium/DebeziumUtils.java  |   0
 .../sort/formats/json/utils/FormatJsonUtil.java    | 165 +++++++++++++++++++++
 .../org.apache.flink.table.factories.Factory       |   0
 .../canal/CanalJsonEnhancedFormatFactoryTest.java  |   0
 .../canal/CanalJsonEnhancedSerDeSchemaTest.java    |   0
 .../json/canal/CanalJsonSerializationTest.java     |   0
 .../json/canal/DebeziumJsonSerializationTest.java  |   0
 .../src/test/resources/canal-json-inlong-data.txt  |   0
 .../src/test/resources/log4j2-test.properties      |   0
 .../sort/formats/json/utils/FormatJsonUtil.java    | 164 ++++++++++++++++++++
 inlong-sort/sort-formats/pom.xml                   |   2 +-
 pom.xml                                            |  24 ++-
 37 files changed, 390 insertions(+), 150 deletions(-)

diff --git a/inlong-manager/manager-plugins/pom.xml 
b/inlong-manager/manager-plugins/pom.xml
index 050addac64..cce8e477e8 100644
--- a/inlong-manager/manager-plugins/pom.xml
+++ b/inlong-manager/manager-plugins/pom.xml
@@ -50,7 +50,7 @@
         </dependency>
         <dependency>
             <groupId>org.apache.inlong</groupId>
-            
<artifactId>sort-flink-dependencies-${sort.flink.version}</artifactId>
+            <artifactId>sort-flink-dependencies-v1.13</artifactId>
             <version>${project.version}</version>
             <exclusions>
                 <exclusion>
diff --git a/inlong-sort/sort-dist/pom.xml b/inlong-sort/sort-dist/pom.xml
index 7f44af9214..9428f1d194 100644
--- a/inlong-sort/sort-dist/pom.xml
+++ b/inlong-sort/sort-dist/pom.xml
@@ -31,8 +31,6 @@
 
     <properties>
         <inlong.root.dir>${project.parent.parent.basedir}</inlong.root.dir>
-        <flink.v15.version>1.15.4</flink.v15.version>
-        <flink.v13.version>1.13.5</flink.v13.version>
     </properties>
 
     <dependencies>
@@ -131,18 +129,18 @@
             <dependencies>
                 <dependency>
                     <groupId>org.apache.inlong</groupId>
-                    <artifactId>sort-format-json</artifactId>
+                    <artifactId>sort-format-json-v1.13</artifactId>
                     <version>${project.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     
<artifactId>flink-sql-parquet_${scala.binary.version}</artifactId>
-                    <version>${flink.v13.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     
<artifactId>flink-sql-orc_${scala.binary.version}</artifactId>
-                    <version>${flink.v13.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
@@ -160,9 +158,6 @@
         </profile>
         <profile>
             <id>v1.15</id>
-            <properties>
-                <flink.version>1.15.4</flink.version>
-            </properties>
             <dependencies>
                 <dependency>
                     <groupId>org.apache.inlong</groupId>
@@ -172,27 +167,27 @@
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     <artifactId>flink-sql-parquet</artifactId>
-                    <version>${flink.v15.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     <artifactId>flink-sql-orc</artifactId>
-                    <version>${flink.v15.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     <artifactId>flink-csv</artifactId>
-                    <version>${flink.v15.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     <artifactId>flink-json</artifactId>
-                    <version>${flink.v15.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
                 <dependency>
                     <groupId>org.apache.flink</groupId>
                     <artifactId>flink-sql-avro</artifactId>
-                    <version>${flink.v15.version}</version>
+                    <version>${flink.version}</version>
                 </dependency>
             </dependencies>
         </profile>
diff --git a/inlong-sort/sort-flink/base/pom.xml 
b/inlong-sort/sort-flink/base/pom.xml
index 57a98050b7..de3eebdc64 100644
--- a/inlong-sort/sort-flink/base/pom.xml
+++ b/inlong-sort/sort-flink/base/pom.xml
@@ -41,11 +41,17 @@
 
         <dependency>
             <groupId>org.apache.inlong</groupId>
-            <artifactId>sort-format-json</artifactId>
+            <artifactId>sort-format-json-${sort.flink.version}</artifactId>
             <version>${project.version}</version>
             <scope>provided</scope>
         </dependency>
 
+        <dependency>
+            <groupId>org.apache.flink</groupId>
+            <artifactId>flink-connector-base</artifactId>
+            <version>${flink.version}</version>
+        </dependency>
+
         <dependency>
             <groupId>org.apache.inlong</groupId>
             <artifactId>sort-format-base</artifactId>
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/sink/log/LogDirtySink.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/sink/log/LogDirtySink.java
index 9bc4125390..2362705f46 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/sink/log/LogDirtySink.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/sink/log/LogDirtySink.java
@@ -21,11 +21,9 @@ import org.apache.inlong.sort.base.dirty.DirtyData;
 import org.apache.inlong.sort.base.dirty.sink.DirtySink;
 import org.apache.inlong.sort.base.dirty.utils.FormatUtils;
 import org.apache.inlong.sort.base.util.LabelUtils;
+import org.apache.inlong.sort.formats.json.utils.FormatJsonUtil;
 
 import org.apache.flink.configuration.Configuration;
-import org.apache.flink.formats.common.TimestampFormat;
-import org.apache.flink.formats.json.JsonOptions.MapNullKeyMode;
-import org.apache.flink.formats.json.RowDataToJsonConverters;
 import 
org.apache.flink.formats.json.RowDataToJsonConverters.RowDataToJsonConverter;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode;
@@ -61,8 +59,7 @@ public class LogDirtySink<T> implements DirtySink<T> {
 
     @Override
     public void open(Configuration configuration) throws Exception {
-        converter = new RowDataToJsonConverters(TimestampFormat.SQL, 
MapNullKeyMode.DROP, null)
-                .createConverter(physicalRowDataType.getLogicalType());
+        converter = FormatJsonUtil.rowDataToJsonConverter(physicalRowDataType);
         fieldGetters = 
FormatUtils.parseFieldGetters(physicalRowDataType.getLogicalType());
     }
 
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/utils/FormatUtils.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/utils/FormatUtils.java
index aca4980db3..110a5dc09c 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/utils/FormatUtils.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/dirty/utils/FormatUtils.java
@@ -17,9 +17,8 @@
 
 package org.apache.inlong.sort.base.dirty.utils;
 
-import org.apache.flink.formats.common.TimestampFormat;
-import org.apache.flink.formats.json.JsonOptions.MapNullKeyMode;
-import org.apache.flink.formats.json.RowDataToJsonConverters;
+import org.apache.inlong.sort.formats.json.utils.FormatJsonUtil;
+
 import 
org.apache.flink.formats.json.RowDataToJsonConverters.RowDataToJsonConverter;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode;
@@ -76,8 +75,7 @@ public final class FormatUtils {
      * @return RowDataToJsonConverter
      */
     public static RowDataToJsonConverter 
parseRowDataToJsonConverter(LogicalType rowType) {
-        return new RowDataToJsonConverters(TimestampFormat.SQL, 
MapNullKeyMode.DROP, null)
-                .createConverter(rowType);
+        return FormatJsonUtil.rowDataToJsonConverter(rowType);
     }
 
     /**
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/DebeziumJsonDynamicSchemaFormat.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/DebeziumJsonDynamicSchemaFormat.java
index 7d996ddade..17d76d9cfc 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/DebeziumJsonDynamicSchemaFormat.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/DebeziumJsonDynamicSchemaFormat.java
@@ -19,20 +19,10 @@ package org.apache.inlong.sort.base.format;
 
 import 
org.apache.inlong.sort.base.format.JsonToRowDataConverters.JsonToRowDataConverter;
 
-import org.apache.flink.shaded.guava18.com.google.common.collect.ImmutableMap;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.types.logical.BigIntType;
-import org.apache.flink.table.types.logical.BooleanType;
-import org.apache.flink.table.types.logical.DoubleType;
-import org.apache.flink.table.types.logical.FloatType;
-import org.apache.flink.table.types.logical.IntType;
 import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.flink.table.types.logical.RowType;
-import org.apache.flink.table.types.logical.SmallIntType;
-import org.apache.flink.table.types.logical.TinyIntType;
-import org.apache.flink.table.types.logical.VarBinaryType;
-import org.apache.flink.table.types.logical.VarCharType;
 import org.apache.flink.types.RowKind;
 
 import java.io.IOException;
@@ -40,6 +30,8 @@ import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
 
+import static 
org.apache.inlong.sort.formats.json.utils.FormatJsonUtil.DEBEZIUM_TYPE_2_FLINK_TYPE_MAPPING;
+
 /**
  * Debezium json dynamic format
  */
@@ -76,19 +68,6 @@ public class DebeziumJsonDynamicSchemaFormat extends 
JsonDynamicSchemaFormat {
      */
     private static final String OP_DELETE = "d";
 
-    private static final Map<String, LogicalType> 
DEBEZIUM_TYPE_2_FLINK_TYPE_MAPPING =
-            ImmutableMap.<String, LogicalType>builder()
-                    .put("BOOLEAN", new BooleanType())
-                    .put("INT8", new TinyIntType())
-                    .put("INT16", new SmallIntType())
-                    .put("INT32", new IntType())
-                    .put("INT64", new BigIntType())
-                    .put("FLOAT32", new FloatType())
-                    .put("FLOAT64", new DoubleType())
-                    .put("STRING", new VarCharType())
-                    .put("BYTES", new VarBinaryType())
-                    .build();
-
     protected DebeziumJsonDynamicSchemaFormat(Map<String, String> props) {
         super(props);
     }
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonDynamicSchemaFormat.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonDynamicSchemaFormat.java
index 6ed59b1de2..51e08f8852 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonDynamicSchemaFormat.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonDynamicSchemaFormat.java
@@ -21,29 +21,13 @@ import org.apache.commons.lang3.StringUtils;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.configuration.ReadableConfig;
 import org.apache.flink.formats.common.TimestampFormat;
-import org.apache.flink.shaded.guava18.com.google.common.collect.ImmutableMap;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.type.TypeReference;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.flink.table.types.logical.BigIntType;
-import org.apache.flink.table.types.logical.BinaryType;
-import org.apache.flink.table.types.logical.BooleanType;
-import org.apache.flink.table.types.logical.CharType;
-import org.apache.flink.table.types.logical.DateType;
 import org.apache.flink.table.types.logical.DecimalType;
-import org.apache.flink.table.types.logical.DoubleType;
-import org.apache.flink.table.types.logical.FloatType;
-import org.apache.flink.table.types.logical.IntType;
-import org.apache.flink.table.types.logical.LocalZonedTimestampType;
 import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.table.types.logical.RowType.RowField;
-import org.apache.flink.table.types.logical.SmallIntType;
-import org.apache.flink.table.types.logical.TimeType;
-import org.apache.flink.table.types.logical.TimestampType;
-import org.apache.flink.table.types.logical.TinyIntType;
-import org.apache.flink.table.types.logical.VarBinaryType;
-import org.apache.flink.table.types.logical.VarCharType;
 import org.apache.flink.types.RowKind;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -58,6 +42,8 @@ import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 
 import static 
org.apache.inlong.sort.base.Constants.SINK_MULTIPLE_TYPE_MAP_COMPATIBLE_WITH_SPARK;
+import static 
org.apache.inlong.sort.formats.json.utils.FormatJsonUtil.SQL_TYPE_2_FLINK_TYPE_MAPPING;
+import static 
org.apache.inlong.sort.formats.json.utils.FormatJsonUtil.SQL_TYPE_2_SPARK_SUPPORTED_FLINK_TYPE_MAPPING;
 
 /**
  * Json dynamic format class
@@ -81,81 +67,11 @@ public abstract class JsonDynamicSchemaFormat extends 
AbstractDynamicSchemaForma
      */
     private static final Integer FIRST = 0;
 
-    private static final int DEFAULT_DECIMAL_PRECISION = 15;
-    private static final int DEFAULT_DECIMAL_SCALE = 5;
     /**
      * dialect sql type pattern such as DECIMAL(38, 10) from mysql or oracle 
etc
      */
     private static final Pattern DIALECT_SQL_TYPE_PATTERN = 
Pattern.compile("\\w+\\(([\\d,\\s]*)\\)");
 
-    private static final Integer ORACLE_TIMESTAMP_TIME_ZONE = -101;
-
-    private static final Map<Integer, LogicalType> 
SQL_TYPE_2_FLINK_TYPE_MAPPING =
-            ImmutableMap.<Integer, LogicalType>builder()
-                    .put(java.sql.Types.CHAR, new CharType())
-                    .put(java.sql.Types.VARCHAR, new VarCharType())
-                    .put(java.sql.Types.SMALLINT, new SmallIntType())
-                    .put(java.sql.Types.INTEGER, new IntType())
-                    .put(java.sql.Types.BIGINT, new BigIntType())
-                    .put(java.sql.Types.REAL, new FloatType())
-                    .put(java.sql.Types.DOUBLE, new DoubleType())
-                    .put(java.sql.Types.FLOAT, new FloatType())
-                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
-                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
-                    .put(java.sql.Types.BIT, new BooleanType())
-                    .put(java.sql.Types.TIME, new TimeType())
-                    .put(java.sql.Types.TIME_WITH_TIMEZONE, new TimeType())
-                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
-                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
-                    .put(java.sql.Types.TIMESTAMP, new TimestampType())
-                    .put(java.sql.Types.BINARY, new BinaryType())
-                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
-                    .put(java.sql.Types.BLOB, new VarBinaryType())
-                    .put(java.sql.Types.CLOB, new VarBinaryType())
-                    .put(java.sql.Types.DATE, new DateType())
-                    .put(java.sql.Types.BOOLEAN, new BooleanType())
-                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
-                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
-                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
-                    .put(java.sql.Types.ARRAY, new VarCharType())
-                    .put(java.sql.Types.NCHAR, new CharType())
-                    .put(java.sql.Types.NCLOB, new VarBinaryType())
-                    .put(java.sql.Types.TINYINT, new TinyIntType())
-                    .put(java.sql.Types.OTHER, new VarCharType())
-                    .build();
-
-    private static final Map<Integer, LogicalType> 
SQL_TYPE_2_SPARK_SUPPORTED_FLINK_TYPE_MAPPING =
-            ImmutableMap.<Integer, LogicalType>builder()
-                    .put(java.sql.Types.CHAR, new CharType())
-                    .put(java.sql.Types.VARCHAR, new VarCharType())
-                    .put(java.sql.Types.SMALLINT, new SmallIntType())
-                    .put(java.sql.Types.INTEGER, new IntType())
-                    .put(java.sql.Types.BIGINT, new BigIntType())
-                    .put(java.sql.Types.REAL, new FloatType())
-                    .put(java.sql.Types.DOUBLE, new DoubleType())
-                    .put(java.sql.Types.FLOAT, new FloatType())
-                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
-                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
-                    .put(java.sql.Types.BIT, new BooleanType())
-                    .put(java.sql.Types.TIME, new VarCharType())
-                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
-                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
-                    .put(java.sql.Types.TIMESTAMP, new 
LocalZonedTimestampType())
-                    .put(java.sql.Types.BINARY, new BinaryType())
-                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
-                    .put(java.sql.Types.BLOB, new VarBinaryType())
-                    .put(java.sql.Types.DATE, new DateType())
-                    .put(java.sql.Types.BOOLEAN, new BooleanType())
-                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
-                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
-                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
-                    .put(java.sql.Types.ARRAY, new VarCharType())
-                    .put(java.sql.Types.NCHAR, new CharType())
-                    .put(java.sql.Types.NCLOB, new VarBinaryType())
-                    .put(java.sql.Types.TINYINT, new TinyIntType())
-                    .put(java.sql.Types.OTHER, new VarCharType())
-                    .build();
-
     public final ObjectMapper objectMapper = new ObjectMapper();
     protected final JsonToRowDataConverters rowDataConverters;
     protected final boolean adaptSparkEngine;
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonToRowDataConverters.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonToRowDataConverters.java
index c0e0d74626..0a77f40aea 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonToRowDataConverters.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/format/JsonToRowDataConverters.java
@@ -39,7 +39,6 @@ import org.apache.flink.table.types.logical.LogicalTypeFamily;
 import org.apache.flink.table.types.logical.MapType;
 import org.apache.flink.table.types.logical.MultisetType;
 import org.apache.flink.table.types.logical.RowType;
-import org.apache.flink.table.types.logical.utils.LogicalTypeChecks;
 import org.apache.flink.table.types.logical.utils.LogicalTypeUtils;
 
 import java.io.IOException;
@@ -328,7 +327,7 @@ public class JsonToRowDataConverters implements 
Serializable {
 
     private JsonToRowDataConverter createMapConverter(
             String typeSummary, LogicalType keyType, LogicalType valueType) {
-        if (!LogicalTypeChecks.hasFamily(keyType, 
LogicalTypeFamily.CHARACTER_STRING)) {
+        if 
(!keyType.getTypeRoot().getFamilies().contains(LogicalTypeFamily.CHARACTER_STRING))
 {
             throw new UnsupportedOperationException(
                     "JSON format doesn't support non-string as key type of 
map. "
                             + "The type is: "
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/sink/MultipleSinkOption.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/sink/MultipleSinkOption.java
index d9086c3cab..821d7ef869 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/sink/MultipleSinkOption.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/sink/MultipleSinkOption.java
@@ -19,11 +19,11 @@ package org.apache.inlong.sort.base.sink;
 
 import org.apache.inlong.sort.schema.TableChange;
 
-import org.apache.flink.shaded.guava18.com.google.common.collect.ImmutableMap;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.Serializable;
+import java.util.HashMap;
 import java.util.Map;
 
 import static 
org.apache.inlong.sort.base.Constants.SINK_MULTIPLE_TYPE_MAP_COMPATIBLE_WITH_SPARK;
@@ -73,8 +73,12 @@ public class MultipleSinkOption implements Serializable {
     }
 
     public Map<String, String> getFormatOption() {
-        return ImmutableMap.of(
-                SINK_MULTIPLE_TYPE_MAP_COMPATIBLE_WITH_SPARK.key(), 
String.valueOf(isSparkEngineEnable()));
+        return new HashMap<String, String>() {
+
+            {
+                put(SINK_MULTIPLE_TYPE_MAP_COMPATIBLE_WITH_SPARK.key(), 
String.valueOf(isSparkEngineEnable()));
+            }
+        };
     }
 
     public SchemaUpdateExceptionPolicy getSchemaUpdatePolicy() {
diff --git 
a/inlong-sort/sort-flink/base/src/test/java/org/apache/inlong/sort/base/dirty/FormatUtilsTest.java
 
b/inlong-sort/sort-flink/base/src/test/java/org/apache/inlong/sort/base/dirty/FormatUtilsTest.java
index 511ad6e215..d748988cdf 100644
--- 
a/inlong-sort/sort-flink/base/src/test/java/org/apache/inlong/sort/base/dirty/FormatUtilsTest.java
+++ 
b/inlong-sort/sort-flink/base/src/test/java/org/apache/inlong/sort/base/dirty/FormatUtilsTest.java
@@ -18,10 +18,8 @@
 package org.apache.inlong.sort.base.dirty;
 
 import org.apache.inlong.sort.base.dirty.utils.FormatUtils;
+import org.apache.inlong.sort.formats.json.utils.FormatJsonUtil;
 
-import org.apache.flink.formats.common.TimestampFormat;
-import org.apache.flink.formats.json.JsonOptions.MapNullKeyMode;
-import org.apache.flink.formats.json.RowDataToJsonConverters;
 import 
org.apache.flink.formats.json.RowDataToJsonConverters.RowDataToJsonConverter;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode;
@@ -87,9 +85,8 @@ public class FormatUtilsTest {
 
     @Test
     public void testJsonFormat() throws JsonProcessingException {
-        RowDataToJsonConverter converter = new 
RowDataToJsonConverters(TimestampFormat.SQL,
-                MapNullKeyMode.DROP, null)
-                        
.createConverter(schema.toPhysicalRowDataType().getLogicalType());
+        RowDataToJsonConverter converter =
+                
FormatJsonUtil.rowDataToJsonConverter(schema.toPhysicalRowDataType().getLogicalType());
         String expectedStr = 
"{\"database\":\"inlong\",\"table\":\"student\",\"id\":1,\"name\":\"leo\",\"age\":18}";
         JsonNode expected = mapper.readValue(expectedStr, JsonNode.class);
         String actual = FormatUtils.jsonFormat(jsonNode, labelMap);
diff --git a/inlong-sort/sort-formats/format-json/pom.xml 
b/inlong-sort/sort-formats/format-json-v1.13/pom.xml
similarity index 96%
rename from inlong-sort/sort-formats/format-json/pom.xml
rename to inlong-sort/sort-formats/format-json-v1.13/pom.xml
index 846219b091..a552e64901 100644
--- a/inlong-sort/sort-formats/format-json/pom.xml
+++ b/inlong-sort/sort-formats/format-json-v1.13/pom.xml
@@ -26,8 +26,8 @@
         <version>1.8.0-SNAPSHOT</version>
     </parent>
 
-    <artifactId>sort-format-json</artifactId>
-    <name>Apache InLong - Sort Format-json</name>
+    <artifactId>sort-format-json-v1.13</artifactId>
+    <name>Apache InLong - Sort Format-json-v1.13</name>
 
     <properties>
         
<inlong.root.dir>${project.parent.parent.parent.basedir}</inlong.root.dir>
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/MysqlBinLogData.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/MysqlBinLogData.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/MysqlBinLogData.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/MysqlBinLogData.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJson.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDecodingFormat.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDecodingFormat.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDecodingFormat.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDecodingFormat.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDeserializationSchema.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDeserializationSchema.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDeserializationSchema.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonDeserializationSchema.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDecodingFormat.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDecodingFormat.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDecodingFormat.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDecodingFormat.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDeserializationSchema.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDeserializationSchema.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDeserializationSchema.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedDeserializationSchema.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedEncodingFormat.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedEncodingFormat.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedEncodingFormat.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedEncodingFormat.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactory.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactory.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactory.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactory.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerializationSchema.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerializationSchema.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerializationSchema.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerializationSchema.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationSchema.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationSchema.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationSchema.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationSchema.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalUtils.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalUtils.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalUtils.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/canal/CanalUtils.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJson.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDecodingFormat.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDecodingFormat.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDecodingFormat.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDecodingFormat.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDeserializationSchema.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDeserializationSchema.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDeserializationSchema.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumJsonDeserializationSchema.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumUtils.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumUtils.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumUtils.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/debezium/DebeziumUtils.java
diff --git 
a/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
new file mode 100644
index 0000000000..cdc15c2324
--- /dev/null
+++ 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
@@ -0,0 +1,165 @@
+/*
+ * 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.sort.formats.json.utils;
+
+import org.apache.flink.formats.common.TimestampFormat;
+import org.apache.flink.formats.json.JsonOptions.MapNullKeyMode;
+import org.apache.flink.formats.json.RowDataToJsonConverters;
+import 
org.apache.flink.formats.json.RowDataToJsonConverters.RowDataToJsonConverter;
+import org.apache.flink.shaded.guava18.com.google.common.collect.ImmutableMap;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.BigIntType;
+import org.apache.flink.table.types.logical.BinaryType;
+import org.apache.flink.table.types.logical.BooleanType;
+import org.apache.flink.table.types.logical.CharType;
+import org.apache.flink.table.types.logical.DateType;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.DoubleType;
+import org.apache.flink.table.types.logical.FloatType;
+import org.apache.flink.table.types.logical.IntType;
+import org.apache.flink.table.types.logical.LocalZonedTimestampType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.SmallIntType;
+import org.apache.flink.table.types.logical.TimeType;
+import org.apache.flink.table.types.logical.TimestampType;
+import org.apache.flink.table.types.logical.TinyIntType;
+import org.apache.flink.table.types.logical.VarBinaryType;
+import org.apache.flink.table.types.logical.VarCharType;
+
+import java.util.Map;
+
+public class FormatJsonUtil {
+
+    private static final int DEFAULT_DECIMAL_PRECISION = 15;
+    private static final int DEFAULT_DECIMAL_SCALE = 5;
+    private static final Integer ORACLE_TIMESTAMP_TIME_ZONE = -101;
+
+    public static RowDataToJsonConverter rowDataToJsonConverter(DataType 
physicalRowDataType) {
+        return rowDataToJsonConverter(TimestampFormat.SQL, null, 
physicalRowDataType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            String mapNullKeyLiteral,
+            DataType physicalRowDataType) {
+        return rowDataToJsonConverter(timestampFormat, MapNullKeyMode.DROP, 
mapNullKeyLiteral, physicalRowDataType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            MapNullKeyMode mapNullKeyMode,
+            String mapNullKeyLiteral, DataType physicalRowDataType) {
+        return new RowDataToJsonConverters(timestampFormat, mapNullKeyMode, 
mapNullKeyLiteral)
+                .createConverter(physicalRowDataType.getLogicalType());
+    }
+
+    public static RowDataToJsonConverter rowDataToJsonConverter(LogicalType 
rowType) {
+        return rowDataToJsonConverter(TimestampFormat.SQL, null, rowType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            String mapNullKeyLiteral,
+            LogicalType rowType) {
+        return rowDataToJsonConverter(timestampFormat, MapNullKeyMode.DROP, 
mapNullKeyLiteral, rowType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            MapNullKeyMode mapNullKeyMode,
+            String mapNullKeyLiteral, LogicalType rowType) {
+        return new RowDataToJsonConverters(timestampFormat, mapNullKeyMode, 
mapNullKeyLiteral)
+                .createConverter(rowType);
+    }
+
+    public static final Map<Integer, LogicalType> 
SQL_TYPE_2_FLINK_TYPE_MAPPING =
+            ImmutableMap.<Integer, LogicalType>builder()
+                    .put(java.sql.Types.CHAR, new CharType())
+                    .put(java.sql.Types.VARCHAR, new VarCharType())
+                    .put(java.sql.Types.SMALLINT, new SmallIntType())
+                    .put(java.sql.Types.INTEGER, new IntType())
+                    .put(java.sql.Types.BIGINT, new BigIntType())
+                    .put(java.sql.Types.REAL, new FloatType())
+                    .put(java.sql.Types.DOUBLE, new DoubleType())
+                    .put(java.sql.Types.FLOAT, new FloatType())
+                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.BIT, new BooleanType())
+                    .put(java.sql.Types.TIME, new TimeType())
+                    .put(java.sql.Types.TIME_WITH_TIMEZONE, new TimeType())
+                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
+                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.TIMESTAMP, new TimestampType())
+                    .put(java.sql.Types.BINARY, new BinaryType())
+                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
+                    .put(java.sql.Types.BLOB, new VarBinaryType())
+                    .put(java.sql.Types.CLOB, new VarBinaryType())
+                    .put(java.sql.Types.DATE, new DateType())
+                    .put(java.sql.Types.BOOLEAN, new BooleanType())
+                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
+                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
+                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
+                    .put(java.sql.Types.ARRAY, new VarCharType())
+                    .put(java.sql.Types.NCHAR, new CharType())
+                    .put(java.sql.Types.NCLOB, new VarBinaryType())
+                    .put(java.sql.Types.TINYINT, new TinyIntType())
+                    .put(java.sql.Types.OTHER, new VarCharType())
+                    .build();
+
+    public static final Map<Integer, LogicalType> 
SQL_TYPE_2_SPARK_SUPPORTED_FLINK_TYPE_MAPPING =
+            ImmutableMap.<Integer, LogicalType>builder()
+                    .put(java.sql.Types.CHAR, new CharType())
+                    .put(java.sql.Types.VARCHAR, new VarCharType())
+                    .put(java.sql.Types.SMALLINT, new SmallIntType())
+                    .put(java.sql.Types.INTEGER, new IntType())
+                    .put(java.sql.Types.BIGINT, new BigIntType())
+                    .put(java.sql.Types.REAL, new FloatType())
+                    .put(java.sql.Types.DOUBLE, new DoubleType())
+                    .put(java.sql.Types.FLOAT, new FloatType())
+                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.BIT, new BooleanType())
+                    .put(java.sql.Types.TIME, new VarCharType())
+                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
+                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.TIMESTAMP, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.BINARY, new BinaryType())
+                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
+                    .put(java.sql.Types.BLOB, new VarBinaryType())
+                    .put(java.sql.Types.DATE, new DateType())
+                    .put(java.sql.Types.BOOLEAN, new BooleanType())
+                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
+                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
+                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
+                    .put(java.sql.Types.ARRAY, new VarCharType())
+                    .put(java.sql.Types.NCHAR, new CharType())
+                    .put(java.sql.Types.NCLOB, new VarBinaryType())
+                    .put(java.sql.Types.TINYINT, new TinyIntType())
+                    .put(java.sql.Types.OTHER, new VarCharType())
+                    .build();
+
+    public static final Map<String, LogicalType> 
DEBEZIUM_TYPE_2_FLINK_TYPE_MAPPING =
+            ImmutableMap.<String, LogicalType>builder()
+                    .put("BOOLEAN", new BooleanType())
+                    .put("INT8", new TinyIntType())
+                    .put("INT16", new SmallIntType())
+                    .put("INT32", new IntType())
+                    .put("INT64", new BigIntType())
+                    .put("FLOAT32", new FloatType())
+                    .put("FLOAT64", new DoubleType())
+                    .put("STRING", new VarCharType())
+                    .put("BYTES", new VarBinaryType())
+                    .build();
+
+}
diff --git 
a/inlong-sort/sort-formats/format-json/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
 
b/inlong-sort/sort-formats/format-json-v1.13/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/main/resources/META-INF/services/org.apache.flink.table.factories.Factory
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactoryTest.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactoryTest.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactoryTest.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedFormatFactoryTest.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerDeSchemaTest.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerDeSchemaTest.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerDeSchemaTest.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonEnhancedSerDeSchemaTest.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationTest.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationTest.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationTest.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/CanalJsonSerializationTest.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/DebeziumJsonSerializationTest.java
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/DebeziumJsonSerializationTest.java
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/java/org/apache/inlong/sort/formats/json/canal/DebeziumJsonSerializationTest.java
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/java/org/apache/inlong/sort/formats/json/canal/DebeziumJsonSerializationTest.java
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/resources/canal-json-inlong-data.txt
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/resources/canal-json-inlong-data.txt
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/resources/canal-json-inlong-data.txt
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/resources/canal-json-inlong-data.txt
diff --git 
a/inlong-sort/sort-formats/format-json/src/test/resources/log4j2-test.properties
 
b/inlong-sort/sort-formats/format-json-v1.13/src/test/resources/log4j2-test.properties
similarity index 100%
rename from 
inlong-sort/sort-formats/format-json/src/test/resources/log4j2-test.properties
rename to 
inlong-sort/sort-formats/format-json-v1.13/src/test/resources/log4j2-test.properties
diff --git 
a/inlong-sort/sort-formats/format-json-v1.15/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
 
b/inlong-sort/sort-formats/format-json-v1.15/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
new file mode 100644
index 0000000000..95db53d218
--- /dev/null
+++ 
b/inlong-sort/sort-formats/format-json-v1.15/src/main/java/org/apache/inlong/sort/formats/json/utils/FormatJsonUtil.java
@@ -0,0 +1,164 @@
+/*
+ * 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.sort.formats.json.utils;
+
+import org.apache.flink.formats.common.TimestampFormat;
+import org.apache.flink.formats.json.JsonFormatOptions.MapNullKeyMode;
+import org.apache.flink.formats.json.RowDataToJsonConverters;
+import 
org.apache.flink.formats.json.RowDataToJsonConverters.RowDataToJsonConverter;
+import org.apache.flink.shaded.guava18.com.google.common.collect.ImmutableMap;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.BigIntType;
+import org.apache.flink.table.types.logical.BinaryType;
+import org.apache.flink.table.types.logical.BooleanType;
+import org.apache.flink.table.types.logical.CharType;
+import org.apache.flink.table.types.logical.DateType;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.DoubleType;
+import org.apache.flink.table.types.logical.FloatType;
+import org.apache.flink.table.types.logical.IntType;
+import org.apache.flink.table.types.logical.LocalZonedTimestampType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.SmallIntType;
+import org.apache.flink.table.types.logical.TimeType;
+import org.apache.flink.table.types.logical.TimestampType;
+import org.apache.flink.table.types.logical.TinyIntType;
+import org.apache.flink.table.types.logical.VarBinaryType;
+import org.apache.flink.table.types.logical.VarCharType;
+
+import java.util.Map;
+
+public class FormatJsonUtil {
+
+    private static final int DEFAULT_DECIMAL_PRECISION = 15;
+    private static final int DEFAULT_DECIMAL_SCALE = 5;
+    private static final Integer ORACLE_TIMESTAMP_TIME_ZONE = -101;
+
+    public static RowDataToJsonConverter rowDataToJsonConverter(DataType 
physicalRowDataType) {
+        return rowDataToJsonConverter(TimestampFormat.SQL, null, 
physicalRowDataType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            String mapNullKeyLiteral,
+            DataType physicalRowDataType) {
+        return rowDataToJsonConverter(timestampFormat, MapNullKeyMode.DROP, 
mapNullKeyLiteral, physicalRowDataType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            MapNullKeyMode mapNullKeyMode,
+            String mapNullKeyLiteral, DataType physicalRowDataType) {
+        return new RowDataToJsonConverters(timestampFormat, mapNullKeyMode, 
mapNullKeyLiteral)
+                .createConverter(physicalRowDataType.getLogicalType());
+    }
+
+    public static RowDataToJsonConverter rowDataToJsonConverter(LogicalType 
rowType) {
+        return rowDataToJsonConverter(TimestampFormat.SQL, null, rowType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            String mapNullKeyLiteral,
+            LogicalType rowType) {
+        return rowDataToJsonConverter(timestampFormat, MapNullKeyMode.DROP, 
mapNullKeyLiteral, rowType);
+    }
+
+    public static RowDataToJsonConverter 
rowDataToJsonConverter(TimestampFormat timestampFormat,
+            MapNullKeyMode mapNullKeyMode,
+            String mapNullKeyLiteral, LogicalType rowType) {
+        return new RowDataToJsonConverters(timestampFormat, mapNullKeyMode, 
mapNullKeyLiteral)
+                .createConverter(rowType);
+    }
+
+    public static final Map<Integer, LogicalType> 
SQL_TYPE_2_FLINK_TYPE_MAPPING =
+            ImmutableMap.<Integer, LogicalType>builder()
+                    .put(java.sql.Types.CHAR, new CharType())
+                    .put(java.sql.Types.VARCHAR, new VarCharType())
+                    .put(java.sql.Types.SMALLINT, new SmallIntType())
+                    .put(java.sql.Types.INTEGER, new IntType())
+                    .put(java.sql.Types.BIGINT, new BigIntType())
+                    .put(java.sql.Types.REAL, new FloatType())
+                    .put(java.sql.Types.DOUBLE, new DoubleType())
+                    .put(java.sql.Types.FLOAT, new FloatType())
+                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.BIT, new BooleanType())
+                    .put(java.sql.Types.TIME, new TimeType())
+                    .put(java.sql.Types.TIME_WITH_TIMEZONE, new TimeType())
+                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
+                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.TIMESTAMP, new TimestampType())
+                    .put(java.sql.Types.BINARY, new BinaryType())
+                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
+                    .put(java.sql.Types.BLOB, new VarBinaryType())
+                    .put(java.sql.Types.CLOB, new VarBinaryType())
+                    .put(java.sql.Types.DATE, new DateType())
+                    .put(java.sql.Types.BOOLEAN, new BooleanType())
+                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
+                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
+                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
+                    .put(java.sql.Types.ARRAY, new VarCharType())
+                    .put(java.sql.Types.NCHAR, new CharType())
+                    .put(java.sql.Types.NCLOB, new VarBinaryType())
+                    .put(java.sql.Types.TINYINT, new TinyIntType())
+                    .put(java.sql.Types.OTHER, new VarCharType())
+                    .build();
+
+    public static final Map<Integer, LogicalType> 
SQL_TYPE_2_SPARK_SUPPORTED_FLINK_TYPE_MAPPING =
+            ImmutableMap.<Integer, LogicalType>builder()
+                    .put(java.sql.Types.CHAR, new CharType())
+                    .put(java.sql.Types.VARCHAR, new VarCharType())
+                    .put(java.sql.Types.SMALLINT, new SmallIntType())
+                    .put(java.sql.Types.INTEGER, new IntType())
+                    .put(java.sql.Types.BIGINT, new BigIntType())
+                    .put(java.sql.Types.REAL, new FloatType())
+                    .put(java.sql.Types.DOUBLE, new DoubleType())
+                    .put(java.sql.Types.FLOAT, new FloatType())
+                    .put(java.sql.Types.DECIMAL, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.NUMERIC, new 
DecimalType(DEFAULT_DECIMAL_PRECISION, DEFAULT_DECIMAL_SCALE))
+                    .put(java.sql.Types.BIT, new BooleanType())
+                    .put(java.sql.Types.TIME, new VarCharType())
+                    .put(java.sql.Types.TIMESTAMP_WITH_TIMEZONE, new 
LocalZonedTimestampType())
+                    .put(ORACLE_TIMESTAMP_TIME_ZONE, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.TIMESTAMP, new 
LocalZonedTimestampType())
+                    .put(java.sql.Types.BINARY, new BinaryType())
+                    .put(java.sql.Types.VARBINARY, new VarBinaryType())
+                    .put(java.sql.Types.BLOB, new VarBinaryType())
+                    .put(java.sql.Types.DATE, new DateType())
+                    .put(java.sql.Types.BOOLEAN, new BooleanType())
+                    .put(java.sql.Types.LONGNVARCHAR, new VarCharType())
+                    .put(java.sql.Types.LONGVARBINARY, new VarCharType())
+                    .put(java.sql.Types.LONGVARCHAR, new VarCharType())
+                    .put(java.sql.Types.ARRAY, new VarCharType())
+                    .put(java.sql.Types.NCHAR, new CharType())
+                    .put(java.sql.Types.NCLOB, new VarBinaryType())
+                    .put(java.sql.Types.TINYINT, new TinyIntType())
+                    .put(java.sql.Types.OTHER, new VarCharType())
+                    .build();
+
+    public static final Map<String, LogicalType> 
DEBEZIUM_TYPE_2_FLINK_TYPE_MAPPING =
+            ImmutableMap.<String, LogicalType>builder()
+                    .put("BOOLEAN", new BooleanType())
+                    .put("INT8", new TinyIntType())
+                    .put("INT16", new SmallIntType())
+                    .put("INT32", new IntType())
+                    .put("INT64", new BigIntType())
+                    .put("FLOAT32", new FloatType())
+                    .put("FLOAT64", new DoubleType())
+                    .put("STRING", new VarCharType())
+                    .put("BYTES", new VarBinaryType())
+                    .build();
+}
diff --git a/inlong-sort/sort-formats/pom.xml b/inlong-sort/sort-formats/pom.xml
index f354a9071f..c698e80f14 100644
--- a/inlong-sort/sort-formats/pom.xml
+++ b/inlong-sort/sort-formats/pom.xml
@@ -37,7 +37,7 @@
         <module>format-kv</module>
         <module>format-inlongmsg-base</module>
         <module>format-inlongmsg-csv</module>
-        <module>format-json</module>
+        <module>format-json-v1.13</module>
         <module>format-inlongmsg-pb</module>
         <module>format-json-v1.15</module>
     </modules>
diff --git a/pom.xml b/pom.xml
index f29de8f2c0..fcd668fb2d 100644
--- a/pom.xml
+++ b/pom.xml
@@ -77,7 +77,6 @@
         <jboss.netty.version>3.10.6.Final</jboss.netty.version>
         <scala.binary.version>2.11</scala.binary.version>
         <spark.version>2.4.4</spark.version>
-        <sort.flink.version>v1.13</sort.flink.version>
 
         
<simpleclient.httpserver.version>0.14.1</simpleclient.httpserver.version>
         <httpcore.version>4.4.14</httpcore.version>
@@ -153,7 +152,8 @@
         <pulsar.version>2.8.1</pulsar.version>
         <kafka.version>2.4.1</kafka.version>
         <iceberg.version>1.1.0</iceberg.version>
-        <flink.version>1.13.5</flink.version>
+        <flink.version.v1.13>1.13.5</flink.version.v1.13>
+        <flink.version.v1.15>1.15.4</flink.version.v1.15>
         <flink.minor.version>1.13</flink.minor.version>
         <flink.scala.binary.version>2.11</flink.scala.binary.version>
         <flink.jackson.version>2.12.1-13.0</flink.jackson.version>
@@ -1510,4 +1510,24 @@
         <url>https://github.com/apache/inlong/actions</url>
     </ciManagement>
 
+    <profiles>
+        <profile>
+            <id>v1.13</id>
+            <activation>
+                <activeByDefault>true</activeByDefault>
+            </activation>
+            <properties>
+                <flink.version>1.13.5</flink.version>
+                <sort.flink.version>v1.13</sort.flink.version>
+            </properties>
+        </profile>
+        <profile>
+            <id>v1.15</id>
+            <properties>
+                <flink.version>1.15.4</flink.version>
+                <sort.flink.version>v1.15</sort.flink.version>
+            </properties>
+        </profile>
+    </profiles>
+
 </project>

Reply via email to