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>