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

zirui 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 e87e23f083 [INLONG-8341][Sort] MySQL cdc connector can get scale for 
decimal field #8342
e87e23f083 is described below

commit e87e23f083e6cb548fa8935226345ea452d837eb
Author: Liao Rui <[email protected]>
AuthorDate: Wed Jun 28 16:16:59 2023 +0800

    [INLONG-8341][Sort] MySQL cdc connector can get scale for decimal field 
#8342
---
 .../inlong/sort/cdc/mysql/utils/MetaDataUtils.java    | 19 +++++++++++++------
 1 file changed, 13 insertions(+), 6 deletions(-)

diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/utils/MetaDataUtils.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/utils/MetaDataUtils.java
index aa4d6ebb9e..4d6c5e8aff 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/utils/MetaDataUtils.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/utils/MetaDataUtils.java
@@ -57,6 +57,9 @@ public class MetaDataUtils {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(MetaDataUtils.class);
 
+    private static final String FORMAT_PRECISION = "%s(%d)";
+    private static final String FORMAT_PRECISION_SCALE = "%s(%d, %d)";
+
     /**
      * get sql type from table schema, represents the jdbc data type
      *
@@ -144,12 +147,16 @@ public class MetaDataUtils {
         table.columns()
                 .forEach(
                         column -> {
-                            mysqlType.put(
-                                    column.name(),
-                                    String.format(
-                                            "%s(%d)",
-                                            column.typeName(),
-                                            column.length()));
+                            if (column.scale().isPresent()) {
+                                mysqlType.put(
+                                        column.name(),
+                                        String.format(FORMAT_PRECISION_SCALE,
+                                                column.typeName(), 
column.length(), column.scale().get()));
+                            } else {
+                                mysqlType.put(
+                                        column.name(),
+                                        String.format(FORMAT_PRECISION, 
column.typeName(), column.length()));
+                            }
                         });
         return mysqlType;
     }

Reply via email to