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,