This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 2833e4855d [Fix][Transform-V2] Fix DECIMAL arithmetic precision loss
and division rounding (#11724)
2833e4855d is described below
commit 2833e4855d228f4a1aaf8f88fc8e4aa9fb751dce
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Wed Aug 19 16:44:45 2026 +0530
[Fix][Transform-V2] Fix DECIMAL arithmetic precision loss and division
rounding (#11724)
---
.../introduction/concepts/incompatible-changes.md | 8 +
.../introduction/concepts/incompatible-changes.md | 8 +
.../transform/sql/zeta/ZetaSQLFunction.java | 58 ++++-
.../transform/sql/SQLDecimalArithmeticTest.java | 274 +++++++++++++++++++++
4 files changed, 340 insertions(+), 8 deletions(-)
diff --git a/docs/en/introduction/concepts/incompatible-changes.md
b/docs/en/introduction/concepts/incompatible-changes.md
index a17d39e087..ae63461640 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -178,6 +178,14 @@ You need to check this document before you upgrade to
related version.
`12345678901234567000.00` and now returns `12345678901234567890.99`, and
`MOD(9007199254740993, 2)` previously
returned `0` and now returns `1`. Jobs that (intentionally or not) depended
on the old lossy values will see
different — now correct — output.
+- **[BREAKING]** SQL Transform arithmetic on `DECIMAL` columns is now exact,
and division rounds to nearest:
+ - Operands of `+`, `-`, `*` and `/` were previously converted with
`BigDecimal.valueOf(value.doubleValue())`, which collapsed them to a `double`
and discarded everything beyond ~17 significant digits. Values now keep full
precision — for example, on `DECIMAL(38,2)` columns `123456789012345678.99 +
0.01` returns `123456789012345679.00` instead of `123456789012345680.01`.
+ - Division now uses `RoundingMode.HALF_UP` instead of `RoundingMode.UP`.
`UP` always rounded away from zero, so at scale 2 `10 / 3` returned `3.34`
instead of `3.33`, and `1 / 1000` returned `0.01` instead of `0.00`.
+ - `%` (`MOD`) is unaffected; it already delegated to the `MOD` function
rather than converting operands itself.
+ - `*` now rounds its result to the scale declared for the output column
(`HALF_UP`), the same way `/` already did. Exact multiplication produces a
result whose scale is the sum of the operand scales, while the column is
declared as `DECIMAL(max(precision), max(scale))`; emitting the wider value
would break sinks that encode against the declared schema. On `DECIMAL(38,2)`
columns `10.25 * 3.75` returns `38.44`, where the old lossy conversion happened
to return `38.4375` for these partic [...]
+ - Dividing by a zero `DECIMAL` now fails with a `TransformException` naming
the operation, where the underlying cause was previously
`java.lang.ArithmeticException("/ by zero")`. The failing expression was
already reported either way, since the SQL engine wraps anything thrown while
evaluating an expression; only the cause type changed. This matches how `MOD`
by zero has always been reported.
+
+ **Migration Guide**: Results that were previously inflated by the old
rounding mode, or truncated by the `double` conversion, will change.
Multiplication results may now carry *fewer* decimal places than before: the
old conversion sometimes emitted a value wider than the declared column scale,
and that value is now rounded down to it, so a job reading `38.4375` from a
`DECIMAL(38,2)` column will read `38.44` after upgrading. Any code that
inspects the *cause* of a division failure and [...]
### Engine Behavior Changes
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md
b/docs/zh/introduction/concepts/incompatible-changes.md
index 359e192edd..fb7373adf3 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -170,6 +170,14 @@
`ROUND(CAST('12345678901234567890.987654321' AS DECIMAL(38,9)), 2)` 以前返回
`12345678901234567000.00`,现在返回
`12345678901234567890.99`;`MOD(9007199254740993, 2)` 以前返回 `0`,
现在返回 `1`。依赖旧的精度丢失结果的作业,其输出会发生变化(现在是正确的)。
+- **[BREAKING]** SQL Transform 对 `DECIMAL` 列的算术运算现在保持精确,并且除法改为四舍五入:
+ - `+`、`-`、`*`、`/` 的操作数之前通过 `BigDecimal.valueOf(value.doubleValue())`
转换,会先退化为 `double`,丢弃约 17 位有效数字之后的全部内容。现在会保留完整精度——例如在 `DECIMAL(38,2)`
列上,`123456789012345678.99 + 0.01` 返回 `123456789012345679.00`,而不是
`123456789012345680.01`。
+ - 除法现在使用 `RoundingMode.HALF_UP` 而不是 `RoundingMode.UP`。`UP` 总是向远离零的方向进位,因此在
scale 为 2 时,`10 / 3` 返回 `3.34` 而不是 `3.33`,`1 / 1000` 返回 `0.01` 而不是 `0.00`。
+ - `%`(`MOD`)不受影响,它本来就委托给 `MOD` 函数,没有自行转换操作数。
+ - `*` 现在会将结果舍入到输出列声明的 scale(`HALF_UP`),与 `/` 的既有行为一致。精确乘法得到的结果 scale 等于两个操作数
scale 之和,而该列声明的类型是 `DECIMAL(max(precision), max(scale))`;若直接输出更宽的值,会导致按声明
schema 编码的 Sink 写入失败。在 `DECIMAL(38,2)` 列上,`10.25 * 3.75` 返回
`38.44`,而旧的有损转换对这组特定的值恰好返回 `38.4375`。
+ - 除数为零的 `DECIMAL` 除法现在抛出标明该运算的 `TransformException`,而此前底层原因是
`java.lang.ArithmeticException("/ by zero")`。两种情况下出错的表达式本来就会被报告(SQL
引擎会包装表达式求值过程中抛出的任何异常),变化的只是 cause 的类型。这与 `MOD` 除零一直以来的报错方式保持一致。
+
+ **迁移指南**:之前被旧舍入模式抬高、或被 `double`
转换截断的结果都会发生变化。乘法结果的小数位数可能比以前*更少*:旧的转换有时会输出比列声明 scale 更宽的值,现在该值会被舍入到声明的
scale,因此原先从 `DECIMAL(38,2)` 列读到 `38.4375` 的作业,升级后会读到
`38.44`。如果下游系统已按旧值对账,升级后需要重新校准。任何为兼容旧行为而做的补偿(例如在除法后减去一个修正值)都应当移除。如果有代码检查除法失败的
cause 并匹配 `ArithmeticException`,需要改为 `TransformException`。
### 引擎行为变更
diff --git
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/sql/zeta/ZetaSQLFunction.java
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/sql/zeta/ZetaSQLFunction.java
index d6a21edc1e..13849dfb27 100644
---
a/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/sql/zeta/ZetaSQLFunction.java
+++
b/seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/sql/zeta/ZetaSQLFunction.java
@@ -761,22 +761,38 @@ public class ZetaSQLFunction {
}
}
if (resultType.getSqlType() == SqlType.DECIMAL) {
- BigDecimal bigDecimal =
BigDecimal.valueOf(leftValue.doubleValue());
+ BigDecimal leftBigDecimal = toBigDecimal(leftValue);
+ BigDecimal rightBigDecimal = toBigDecimal(rightValue);
if (binaryExpression instanceof Addition) {
- return
bigDecimal.add(BigDecimal.valueOf(rightValue.doubleValue()));
+ return leftBigDecimal.add(rightBigDecimal);
}
if (binaryExpression instanceof Subtraction) {
- return
bigDecimal.subtract(BigDecimal.valueOf(rightValue.doubleValue()));
+ return leftBigDecimal.subtract(rightBigDecimal);
}
if (binaryExpression instanceof Multiplication) {
- return
bigDecimal.multiply(BigDecimal.valueOf(rightValue.doubleValue()));
+ // BigDecimal.multiply returns a value whose scale is
leftScale + rightScale,
+ // but the column type declared for this expression is
+ // DECIMAL(max(precision), max(scale)). Normalise to the
declared scale, as
+ // Division already does, so the emitted value matches the
schema the transform
+ // advertises: a sink that encodes against that schema
(Parquet through Avro,
+ // for example) rejects a value whose scale differs from the
declared one.
+ DecimalType decimalType = (DecimalType) resultType;
+ return leftBigDecimal
+ .multiply(rightBigDecimal)
+ .setScale(decimalType.getScale(),
RoundingMode.HALF_UP);
}
if (binaryExpression instanceof Division) {
+ if (rightBigDecimal.signum() == 0) {
+ // BigDecimal.divide would throw a bare
ArithmeticException("/ by zero") with
+ // no indication of which expression produced it. MOD
already reports this as
+ // a TransformException; do the same here, and name the
expression.
+ throw new TransformException(
+ CommonErrorCodeDeprecated.UNSUPPORTED_OPERATION,
+ String.format("Division by zero in expression:
%s", binaryExpression));
+ }
DecimalType decimalType = (DecimalType) resultType;
- return bigDecimal.divide(
- BigDecimal.valueOf(rightValue.doubleValue()),
- decimalType.getScale(),
- RoundingMode.UP);
+ return leftBigDecimal.divide(
+ rightBigDecimal, decimalType.getScale(),
RoundingMode.HALF_UP);
}
if (binaryExpression instanceof Modulo) {
List<Object> args = new ArrayList<>();
@@ -824,6 +840,32 @@ public class ZetaSQLFunction {
String.format("Unsupported SQL Expression: %s ",
binaryExpression));
}
+ /**
+ * Converts a numeric operand of a DECIMAL expression to {@link
BigDecimal} without routing it
+ * through {@code double}.
+ *
+ * <p>{@code BigDecimal.valueOf(value.doubleValue())} would collapse the
operand to a {@code
+ * double} first, discarding everything beyond ~17 significant digits
before the arithmetic even
+ * starts, which defeats the purpose of the DECIMAL type.
+ *
+ * @param value operand of a binary DECIMAL expression
+ * @return the operand as an exact BigDecimal
+ */
+ private static BigDecimal toBigDecimal(Number value) {
+ if (value instanceof BigDecimal) {
+ return (BigDecimal) value;
+ }
+ if (value instanceof Byte
+ || value instanceof Short
+ || value instanceof Integer
+ || value instanceof Long) {
+ return BigDecimal.valueOf(value.longValue());
+ }
+ // Float/Double have no exact decimal form; valueOf uses the canonical
shortest
+ // representation, which is the closest thing to the value the user
wrote.
+ return BigDecimal.valueOf(value.doubleValue());
+ }
+
public List<SeaTunnelRow> lateralView(
List<SeaTunnelRow> seaTunnelRows,
List<LateralView> lateralViews,
diff --git
a/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/sql/SQLDecimalArithmeticTest.java
b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/sql/SQLDecimalArithmeticTest.java
new file mode 100644
index 0000000000..0705ce59e1
--- /dev/null
+++
b/seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/sql/SQLDecimalArithmeticTest.java
@@ -0,0 +1,274 @@
+/*
+ * 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.seatunnel.transform.sql;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.DecimalType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.transform.exception.TransformException;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.math.BigDecimal;
+import java.util.Collections;
+import java.util.List;
+
+/** Tests that binary arithmetic on DECIMAL columns stays exact and rounds
division to nearest. */
+public class SQLDecimalArithmeticTest {
+
+ private SeaTunnelRow runSql(String query, SeaTunnelRowType rowType,
Object... values) {
+ CatalogTable table = CatalogTableUtil.getCatalogTable("test", rowType);
+ ReadonlyConfig config =
ReadonlyConfig.fromMap(Collections.singletonMap("query", query));
+ SQLTransform transform = new SQLTransform(config, table);
+ List<SeaTunnelRow> out = transform.transformRow(new
SeaTunnelRow(values));
+ Assertions.assertNotNull(out);
+ Assertions.assertFalse(out.isEmpty());
+ return out.get(0);
+ }
+
+ private static String plain(Object field) {
+ Assertions.assertInstanceOf(BigDecimal.class, field);
+ return ((BigDecimal) field).toPlainString();
+ }
+
+ private static SeaTunnelRowType twoDecimals(int precision, int scale) {
+ return new SeaTunnelRowType(
+ new String[] {"a", "b"},
+ new SeaTunnelDataType[] {
+ new DecimalType(precision, scale), new
DecimalType(precision, scale)
+ });
+ }
+
+ /**
+ * A DECIMAL(38,2) value with 20 significant digits does not survive a
round trip through
+ * double, so the operands must be converted exactly. The name says
"operands" rather than
+ * "results": {@code +} and {@code -} are exact end to end, but {@code *}
is computed exactly
+ * and then rounded to the scale declared for its column, so its result is
exact only up to that
+ * scale.
+ */
+ @Test
+ public void testAddSubtractMultiplyUseExactOperands() {
+ SeaTunnelRowType rowType = twoDecimals(38, 2);
+
+ SeaTunnelRow outRow =
+ runSql(
+ "select a + b as sum_val,"
+ + " a - b as diff_val,"
+ + " a * b as mul_val"
+ + " from dual",
+ rowType,
+ new BigDecimal("123456789012345678.99"),
+ new BigDecimal("0.01"));
+
+ // Before the fix these were 123456789012345680.01,
123456789012345680.00 and
+ // 1234567890123456.8 respectively.
+ Assertions.assertEquals("123456789012345679.00",
plain(outRow.getField(0)));
+ Assertions.assertEquals("123456789012345678.98",
plain(outRow.getField(1)));
+ // The exact product is 1234567890123456.7899; it is rounded to the
scale declared for
+ // the output column, the same way division is.
+ Assertions.assertEquals("1234567890123456.79",
plain(outRow.getField(2)));
+ }
+
+ /** An integral operand wider than double's 53-bit mantissa must not be
rounded either. */
+ @Test
+ public void testMixedDecimalAndBigintStaysExact() {
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"a", "b"},
+ new SeaTunnelDataType[] {new DecimalType(38, 0),
BasicType.LONG_TYPE});
+
+ SeaTunnelRow outRow =
+ runSql(
+ "select a + b as sum_val from dual",
+ rowType,
+ BigDecimal.ONE,
+ 9007199254740993L);
+
+ // 9007199254740993 is not representable as a double; it collapses to
...992, so before
+ // the fix this returned 9007199254740993 instead of 9007199254740994.
+ Assertions.assertEquals("9007199254740994", plain(outRow.getField(0)));
+ }
+
+ /**
+ * Division must round to nearest. RoundingMode.UP always rounds away from
zero, which inflates
+ * every inexact quotient.
+ */
+ @Test
+ public void testDivisionRoundsHalfUp() {
+ SeaTunnelRow outRow =
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("10.00"),
+ new BigDecimal("3.00"));
+
+ // RoundingMode.UP gave 3.34.
+ Assertions.assertEquals("3.33", plain(outRow.getField(0)));
+ }
+
+ /** A quotient below half the last representable digit must round down to
zero, not up to it. */
+ @Test
+ public void testDivisionDoesNotInventValue() {
+ SeaTunnelRow outRow =
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("1.00"),
+ new BigDecimal("1000.00"));
+
+ // RoundingMode.UP gave 0.01, manufacturing value out of a quotient
that rounds to zero.
+ Assertions.assertEquals("0.00", plain(outRow.getField(0)));
+ }
+
+ /**
+ * Negative quotients must round to nearest as well. RoundingMode.UP
rounds away from zero in
+ * both directions, so it made negative results more negative rather than
closer to zero.
+ */
+ @Test
+ public void testNegativeDivisionRoundsHalfUp() {
+ SeaTunnelRow inexact =
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("-10.00"),
+ new BigDecimal("3.00"));
+ // RoundingMode.UP gave -3.34.
+ Assertions.assertEquals("-3.33", plain(inexact.getField(0)));
+
+ SeaTunnelRow towardsZero =
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("-1.00"),
+ new BigDecimal("1000.00"));
+ // RoundingMode.UP gave -0.01, inventing a debit out of a quotient
that rounds to zero.
+ Assertions.assertEquals("0.00", plain(towardsZero.getField(0)));
+ }
+
+ /**
+ * HALF_UP breaks an exact tie by rounding away from zero, so -0.005 at
scale 2 is -0.01 rather
+ * than 0.00. Pinned here because it is the one case where HALF_UP and the
old UP agree, and a
+ * future switch to HALF_EVEN would silently change it.
+ */
+ @Test
+ public void testNegativeDivisionTieRoundsAwayFromZero() {
+ SeaTunnelRow outRow =
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("-5.00"),
+ new BigDecimal("1000.00"));
+
+ Assertions.assertEquals("-0.01", plain(outRow.getField(0)));
+ }
+
+ /**
+ * Dividing by a zero DECIMAL surfaces BigDecimal's bare
ArithmeticException("/ by zero") as the
+ * cause. It is now reported the same way MOD by zero already is, naming
the operation.
+ */
+ @Test
+ public void testDivisionByZeroReportsExpression() {
+ TransformException exception =
+ Assertions.assertThrows(
+ TransformException.class,
+ () ->
+ runSql(
+ "select a / b as div_val from dual",
+ twoDecimals(38, 2),
+ new BigDecimal("1.00"),
+ new BigDecimal("0.00")));
+
+ // ZetaSQLEngine already wraps any failure with the expression that
produced it.
+ Assertions.assertTrue(exception.getMessage().contains("a / b"),
exception.getMessage());
+
+ // The cause is what changes here: previously ArithmeticException("/
by zero").
+ Throwable cause = exception.getCause();
+ Assertions.assertInstanceOf(TransformException.class, cause);
+ Assertions.assertTrue(cause.getMessage().contains("Division by zero"),
cause.getMessage());
+ }
+
+ /**
+ * When both operands are DECIMAL, every emitted DECIMAL must carry the
scale that the transform
+ * declares for its column. A sink that builds its write schema from the
declared type and then
+ * encodes the value against it rejects the row when the two disagree, so
an exact result is not
+ * usable on its own.
+ *
+ * <p>The invariant is asserted only for DECIMAL-on-DECIMAL arithmetic,
which is what this test
+ * exercises. It does not hold in general: {@code getExpressionType}
declares DECIMAL as soon as
+ * either side is DECIMAL, while {@code BigDecimal.add}/{@code subtract}
return the max of the
+ * operands' <em>runtime</em> scales, so a FLOAT/DOUBLE operand can still
push {@code +} and
+ * {@code -} past the declared scale. That path is unchanged from before
this fix and is tracked
+ * separately.
+ *
+ * <p>Precision is asserted alongside scale because {@code setScale}
bounds only the latter. The
+ * check documents the range these operands stay within; it is not a proof
that the declared
+ * precision can never be exceeded, since nothing in the DECIMAL branch
bounds it.
+ */
+ @Test
+ public void testEmittedScaleMatchesDeclaredType() {
+ SeaTunnelRowType rowType = twoDecimals(38, 2);
+ CatalogTable table = CatalogTableUtil.getCatalogTable("test", rowType);
+ ReadonlyConfig config =
+ ReadonlyConfig.fromMap(
+ Collections.singletonMap(
+ "query",
+ "select a + b as sum_val,"
+ + " a - b as diff_val,"
+ + " a * b as mul_val,"
+ + " a / b as div_val"
+ + " from dual"));
+ SQLTransform transform = new SQLTransform(config, table);
+
+ SeaTunnelRowType outType =
transform.getProducedCatalogTable().getSeaTunnelRowType();
+ SeaTunnelRow outRow =
+ transform
+ .transformRow(
+ new SeaTunnelRow(
+ new Object[] {
+ new BigDecimal("10.25"), new
BigDecimal("3.75")
+ }))
+ .get(0);
+
+ for (int i = 0; i < outType.getTotalFields(); i++) {
+ SeaTunnelDataType<?> fieldType = outType.getFieldType(i);
+ Assertions.assertInstanceOf(DecimalType.class, fieldType);
+ Assertions.assertEquals(
+ ((DecimalType) fieldType).getScale(),
+ ((BigDecimal) outRow.getField(i)).scale(),
+ "declared and emitted scale differ for column "
+ + outType.getFieldName(i)
+ + " (value "
+ + outRow.getField(i)
+ + ")");
+ Assertions.assertTrue(
+ ((BigDecimal) outRow.getField(i)).precision()
+ <= ((DecimalType) fieldType).getPrecision(),
+ "emitted precision exceeds the declared precision for
column "
+ + outType.getFieldName(i)
+ + " (value "
+ + outRow.getField(i)
+ + ")");
+ }
+ }
+}