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

gongzhongqiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git


The following commit(s) were added to refs/heads/master by this push:
     new c5f391cf6 [FLINK-35865][base] Support Byte and Short in ObjectUtils 
(#3481)
c5f391cf6 is described below

commit c5f391cf6a032333f398038cf96b11a0354bf3c3
Author: gongzhongqiang <[email protected]>
AuthorDate: Tue Jul 23 10:46:18 2024 +0800

    [FLINK-35865][base] Support Byte and Short in ObjectUtils (#3481)
---
 .../cdc/composer/flink/FlinkEnvironmentUtils.java  |  2 +-
 .../cdc/connectors/base/utils/ObjectUtils.java     | 22 ++++++++++++++++++++--
 .../assigner/splitter/OracleChunkSplitter.java     |  6 ++++--
 3 files changed, 25 insertions(+), 5 deletions(-)

diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkEnvironmentUtils.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkEnvironmentUtils.java
index b00717cb0..00adb1e4b 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkEnvironmentUtils.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkEnvironmentUtils.java
@@ -65,7 +65,7 @@ public class FlinkEnvironmentUtils {
                     Stream.concat(previousJars.stream(), 
jarUrls.stream().map(URL::toString))
                             .distinct()
                             .collect(Collectors.toList());
-            LOG.info("pipeline.jars is " + String.join(",", currentJars));
+            LOG.info("pipeline.jars is {}", String.join(",", currentJars));
             configuration.set(PipelineOptions.JARS, currentJars);
         } catch (Exception e) {
             throw new RuntimeException("Failed to add JAR to Flink execution 
environment", e);
diff --git 
a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/utils/ObjectUtils.java
 
b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/utils/ObjectUtils.java
index ea4896fe9..74b3d3af5 100644
--- 
a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/utils/ObjectUtils.java
+++ 
b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/utils/ObjectUtils.java
@@ -28,7 +28,19 @@ public class ObjectUtils {
      * will throw {@link ArithmeticException} if number overflows.
      */
     public static Object plus(Object number, int augend) throws 
ArithmeticException {
-        if (number instanceof Integer) {
+        if (number instanceof Byte) {
+            int result = Math.addExact((Byte) number, augend);
+            if (result < Byte.MIN_VALUE || result > Byte.MAX_VALUE) {
+                throw new ArithmeticException("byte overflow");
+            }
+            return (byte) result;
+        } else if (number instanceof Short) {
+            int result = Math.addExact((Short) number, augend);
+            if (result < Short.MIN_VALUE || result > Short.MAX_VALUE) {
+                throw new ArithmeticException("short overflow");
+            }
+            return (short) result;
+        } else if (number instanceof Integer) {
             return Math.addExact((Integer) number, augend);
         } else if (number instanceof Long) {
             return Math.addExact((Long) number, augend);
@@ -53,7 +65,13 @@ public class ObjectUtils {
                             minuend.getClass().getSimpleName(),
                             subtrahend.getClass().getSimpleName()));
         }
-        if (minuend instanceof Integer) {
+        if (minuend instanceof Byte) {
+            return BigDecimal.valueOf((byte) minuend)
+                    .subtract(BigDecimal.valueOf((byte) subtrahend));
+        } else if (minuend instanceof Short) {
+            return BigDecimal.valueOf((short) minuend)
+                    .subtract(BigDecimal.valueOf((short) subtrahend));
+        } else if (minuend instanceof Integer) {
             return BigDecimal.valueOf((int) 
minuend).subtract(BigDecimal.valueOf((int) subtrahend));
         } else if (minuend instanceof Long) {
             return BigDecimal.valueOf((long) minuend)
diff --git 
a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/assigner/splitter/OracleChunkSplitter.java
 
b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/assigner/splitter/OracleChunkSplitter.java
index 6dfb16e0b..d68f1cd99 100644
--- 
a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/assigner/splitter/OracleChunkSplitter.java
+++ 
b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-oracle-cdc/src/main/java/org/apache/flink/cdc/connectors/oracle/source/assigner/splitter/OracleChunkSplitter.java
@@ -40,6 +40,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.math.BigDecimal;
+import java.math.RoundingMode;
 import java.sql.SQLException;
 import java.util.ArrayList;
 import java.util.Collection;
@@ -49,7 +50,6 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 
-import static java.math.BigDecimal.ROUND_CEILING;
 import static 
org.apache.flink.cdc.connectors.base.utils.ObjectUtils.doubleCompare;
 
 /**
@@ -351,7 +351,9 @@ public class OracleChunkSplitter implements 
JdbcSourceChunkSplitter {
         // factor = (max - min + 1) / rowCount
         final BigDecimal subRowCnt = difference.add(BigDecimal.valueOf(1));
         double distributionFactor =
-                subRowCnt.divide(new BigDecimal(approximateRowCnt), 4, 
ROUND_CEILING).doubleValue();
+                subRowCnt
+                        .divide(new BigDecimal(approximateRowCnt), 4, 
RoundingMode.CEILING)
+                        .doubleValue();
         LOG.info(
                 "The distribution factor of table {} is {} according to the 
min split key {}, max split key {} and approximate row count {}",
                 tableId,

Reply via email to