[SPARK-59805][SQL] Avoid BigDecimal round trips for wide decimals that fit in a long in UnsafeRow - #59159
Open
dwsmith1983 wants to merge 1 commit into
Open
[SPARK-59805][SQL] Avoid BigDecimal round trips for wide decimals that fit in a long in UnsafeRow#59159dwsmith1983 wants to merge 1 commit into
dwsmith1983 wants to merge 1 commit into
Conversation
…t fit in a long in UnsafeRow For DECIMAL columns with precision > 18, UnsafeRow.getDecimal, UnsafeRow.setDecimal and UnsafeRowWriter.write always went through BigInteger/BigDecimal, even when the unscaled value fits in a long. Decode and encode such values directly as a long (the stored bytes stay identical to BigInteger.toByteArray) and return a compact Decimal when the unscaled value is below 10^18 and the scale is between 0 and 18. This lets sum/avg over e.g. decimal(15,2), whose buffer is decimal(25,2), stay on Decimal's existing long fast path. Also make Decimal.roundToInt/Short/Byte truncate exactly via BigInteger for BigDecimal-backed values, as roundToLong already does, instead of going through a Double that rounds values just below a bound (e.g. 2147483647.99999999) past it and wrongly overflows.
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
For
DECIMALcolumns with precision > 18,UnsafeRowstores the unscaled value as big-endian two's-complement bytes. Reading and writing them always went throughBigInteger/BigDecimal, even when the value fits in along:UnsafeRow.getDecimalallocated abyte[], aBigInteger, ajava.math.BigDecimaland a ScalaBigDecimalper call, and returned a BigDecimal-backedDecimal.UnsafeRow.setDecimalandUnsafeRowWriter.writecalledtoJavaBigDecimal().unscaledValue().toByteArray().This PR adds a fast path for values whose unscaled value fits in a
long:UnsafeRow.getDecimal: when the stored value is at most 8 bytes, decode it directly from the row into along. If|unscaled| < 10^18and0 <= scale <= 18, return a compact (long-backed)DecimalviaDecimal.createUnsafe; otherwise return the same BigDecimal-backedDecimalas before (built withBigDecimal.valueOf, without thebyte[]/BigInteger). Negative scales, possible only withspark.sql.legacy.allowNegativeScaleOfDecimal, keep the BigDecimal-backed result.UnsafeRow.setDecimal/UnsafeRowWriter.write: after the existingchangePrecision, a compactDecimalis written with the newUnsafeRow.writeCompactUnscaledBytes, which emits the minimal big-endian two's-complement bytes of thelong. These are byte-identical toBigInteger.toByteArray(), so the row format, row hashing and row equality are unchanged. Non-compact values keep the existing path.Decimal.roundToInt/roundToShort/roundToBytenow truncate exactly viaBigIntegerfor BigDecimal-backed values, asroundToLongalready does, instead of going through aDouble.The main effect is on aggregation: the buffer of
sum/avgoverdecimal(p, s)has precisionp + 10, so e.g.sum(decimal(15,2))uses adecimal(25,2)buffer. Previously every buffer read came back BigDecimal-backed, soDecimal.+never took its existing long fast path; now it does while the running value fits in a long.Why are the changes needed?
Performance. The per-row
BigInteger/BigDecimalround trips are a large share of CPU in decimal-heavy aggregations. In a CPU profile of TPC-H SF10, 11.4% of samples were injava.math.BigDecimal/BigIntegerorDecimal.Item 3 is needed for consistency. With the double-based check, a value just below a bound, such as
2147483647.99999999asdecimal(20,8), rounds to2147483648.0andCAST(... AS INT)wrongly fails withCAST_OVERFLOWunder ANSI mode. Long-backed values of the same number were already checked exactly. Because this PR returns more wide decimals in the compact form, leaving the double-based check would make the result of the same cast depend on the internal representation.Does this PR introduce any user-facing change?
Yes, a bug fix. Casting a
DECIMALwhose value truncates to exactlyInt.MaxValue/Int.MinValue(and likewise forSMALLINT/TINYINT), but which is close enough to the next integer to round up past the bound as a double, no longer fails withCAST_OVERFLOWunder ANSI mode (or returns null undertry_cast). It returns the truncated value, as it already did for the same value held as an unscaled long. For example,CAST(CAST('2147483647.99999999' AS DECIMAL(20,8)) AS INT)now returns2147483647instead of failing.No other change: query results, the
UnsafeRowbinary format and plans are unchanged.How was this patch tested?
New tests:
UnsafeRowWriterSuite: fixed-seed round trips for precision 19..38 across scales. Values include 0, +-1, +-(10^18 - 1), +-10^18, every byte-length boundary +-2^(8k-1) and its neighbours (includingLong.MinValue/Long.MaxValue), values over 8 bytes, and random values. Each is written throughUnsafeRowWriterand through in-placesetDecimal. The tests assert:BigInteger.toByteArray;getDecimalmatches the previous BigDecimal-based result in value, precision, scale,equals,hashCodeandtoString;toLong,toInt,roundToLong,roundToInt,roundToShort,roundToByte,floorandceilgive the same value, or throw the same exception;|unscaled| < 10^18and0 <= scale <= 18.UnsafeRowWriterSuitealso covers rescaling in the writer, overflow/null handling (null bit set, slot kept, re-update after null), and negative scales underspark.sql.legacy.allowNegativeScaleOfDecimal.CastWithAnsiOnSuite(inherited byTryCastSuite): casts of BigDecimal-backeddecimal(20, s)values at theINT/SMALLINT/TINYINTbounds return the bound, and one unit in the last place beyond each bound fails withCAST_OVERFLOW(null undertry_cast), interpreted and codegen.DecimalSuite:roundToInt/roundToShort/roundToByteon BigDecimal-backed values just inside and just outside each bound, compared with the long-backed representation.DataFrameAggregateSuite:sum/avgofdecimal(15,2), grouped and global, with ANSI on and off. Within a single partition, running sums cross unscaled +-10^18 and come back below it, in both directions; other groups cross only when partial aggregates are merged.Existing suites pass:
DecimalSuite,UnsafeRowWriterSuite,UnsafeRowConverterSuite,DecimalExpressionSuite,CastWithAnsiOnSuite,CastWithAnsiOffSuite,TryCastSuiteandDataFrameAggregateSuite. No golden files changed.Benchmark: a new
DecimalArithmeticBenchmark(5M rows, 1000 groups; ns per row; ANSI off / on):The q1 expression doesn't improve because the multiply still produces BigDecimal-backed values; that is left for a follow-up. The
*-results.txtfiles are not included; they should be generated with the GitHub Actions benchmark workflow.End to end, TPC-H SF10 (local[4], median of 3 runs) takes 73.8s instead of 79.4s; q1 goes from 11.55s to 8.76s. TPC-DS SF10 is unchanged within noise (its decimals are
decimal(7,2), whose sums fit in 18 digits).Was this patch authored or co-authored using generative AI tooling?
No