You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
[EPIC] Faster Spark-to-Arrow conversion in CometSparkToColumnarExec #6565
CometSparkToColumnarExec converts Spark rows and Spark column vectors to Arrow so that native operators can consume them. The conversion goes through ArrowWriter, which CometLocalTableScanExec and the Comet cache serializer also use. I measured it one column type at a time on main at 93d1189be, with 8192-row batches. Several common cases are far slower than they need to be:
Decimals cost 30-160 ns per row, 30 to 100 times the cost of an int. Arrow's BigDecimal setter allocates a BigInteger and two byte arrays per value, on top of the Decimal and BigDecimal that Spark's getters build.
The bulk copy for fixed-width columns (perf: bulk copy fixed-width columns in ArrowWriter #5442) is skipped when a batch holds a single null, and for every dictionary-encoded column. Spark's Parquet reader leaves low-cardinality columns dictionary-encoded.
Strings cost 13-16 ns per row, written one value at a time through Arrow's setSafe. A dictionary-encoded string also allocates a copy of its bytes on every row, because Parquet's dictionary returns a fresh array for each decode.
Arrays, maps and structs are written one element at a time: about 90 ns per row for array<string> and 200 ns for map<string,string>, with five entries on average.
Spark's vectorized Parquet reader produces 4096-row batches, and each becomes one Arrow batch, so native operators get half of Comet's batch size. The cache produces 10000-row batches, which become alternating Arrow batches of 8192 and 1808 rows.
End to end, conversion took 549 ms of a 1088 ms TPC-H Q1-style query over Spark's Parquet reader.
What / Why
CometSparkToColumnarExecconverts Spark rows and Spark column vectors to Arrow so that native operators can consume them. The conversion goes throughArrowWriter, whichCometLocalTableScanExecand the Comet cache serializer also use. I measured it one column type at a time onmainat93d1189be, with 8192-row batches. Several common cases are far slower than they need to be:BigDecimalsetter allocates aBigIntegerand two byte arrays per value, on top of theDecimalandBigDecimalthat Spark's getters build.setSafe. A dictionary-encoded string also allocates a copy of its bytes on every row, because Parquet's dictionary returns a fresh array for each decode.array<string>and 200 ns formap<string,string>, with five entries on average.End to end, conversion took 549 ms of a 1088 ms TPC-H Q1-style query over Spark's Parquet reader.
Work
allocateNew) is about half the remaining cost of the fixed-width columnar path once the above lands.CometSparkToColumnarExec.isTypeSupportedstill declines arrays other thanARRAY<STRING>and maps other thanMAP<STRING,STRING>, so those columns never reach the conversion and the operators above them stay on Spark. In the feat: admit string maps in Spark-to-Comet conversion #6036 review, dropping that gate passed every probe query.conversionTimemetric includes the time the child spends producing rows, so it overstates the conversion.SparkColumnarArrowReaderends each Arrow batch at an input file boundary. That protects nothing once fix: keep a Spark scan unconverted when the plan reads input_file_name #6703 makes the conversion fall back forinput_file_name, and it leaves small batches when a task reads many small files.