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)
+                            + ")");
+        }
+    }
+}

Reply via email to