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
#5868 calls Arrow's unshred_variant, then decodes and rebuilds the result in rebuild_spark_variant. Emit Spark-compatible bytes during the first traversal to remove this repeated work. This is the first whole-value performance priority.
Proposed change
Extend Arrow's recursive unshredder with an optional Spark output writer, reusing its typed decoders and value/object/list builders. The public row-decoder prerequisite is tracked by apache/arrow-rs#11260; it is proposed and not available in Comet's pinned dependency.
Normalize typed scalars at their append sites: value-dependent integer/decimal widths, short strings and canonical Java NaN bits, preserving signed zero. Preserve residual scalar encodings, including wide integers/strings and NaN payloads.
Build a fresh dictionary per shredded row. Register emitted typed names in schema order, recurse before the next sibling, then visit residual fields in their original order. Resolve residual names through the input dictionary and remap IDs in objects and lists. [Variant] Replace missing-key metadata repair with native permissive reconstruction #5979 owns this metadata path.
Let ObjectBuilder use a metadata-supplied comparator: UTF-8 by default, the supported Spark profile's ordering in Spark mode. Sort field references without changing dictionary or payload insertion order. Emit Spark's metadata flags.
Read legacy residuals with checked, fallible iteration that preserves their original field sequence. Validate bounds, IDs, types, nesting, unique names and canonical/legacy ordering while descending. Sorting a temporary residual first changes Spark's dictionary IDs. Keep the existing legacy path until this reader is available.
The Spark reconstruction and builder contract was checked on 4.0.4, 4.1.3 and 4.2.0. Those profiles use UTF-16 object-key ordering and clear the metadata sorted flag. #5474 tracks changing ordering when supported Spark versions permit it.
After integration, remove rebuild_spark_variant and its replaced local writer helpers. A final BinaryView-to-Binary layout conversion is acceptable. Retain storage normalization, empty-key repair and other preparation until their own replacements land (#5477); Arrow/DataFusion upgrades are currently deferred.
Completion and performance
Keep exact value and metadata comparisons for scalar-width boundaries, 63/64-byte strings, typed/residual NaNs, signed zero, dictionary traversal and ID/offset widths, Unicode and empty keys, nested residual containers, unused keys and parent nulls. Run native scan tests and unchanged upstream Spark assertions for all supported profiles.
Compare the merged baseline, candidate and vanilla Spark on matched canonical, partially shredded, fully shredded and empty-key files. Consume both byte columns, run both reader orders, and report allocation traffic separately from scan time. The target is faster whole-value reads than vanilla Spark across these inputs; keep benchmark results in the implementing PR description.
Incremental optimizations may land first. Close this issue only when the extra reconstruction pass is removed for all currently supported inputs and the compatibility and performance targets are met.
#5868 calls Arrow's
unshred_variant, then decodes and rebuilds the result inrebuild_spark_variant. Emit Spark-compatible bytes during the first traversal to remove this repeated work. This is the first whole-value performance priority.Proposed change
Extend Arrow's recursive unshredder with an optional Spark output writer, reusing its typed decoders and value/object/list builders. The public row-decoder prerequisite is tracked by apache/arrow-rs#11260; it is proposed and not available in Comet's pinned dependency.
ObjectBuilderuse a metadata-supplied comparator: UTF-8 by default, the supported Spark profile's ordering in Spark mode. Sort field references without changing dictionary or payload insertion order. Emit Spark's metadata flags.typed_value. A NULL typed value in a shredded schema still requires a fresh output dictionary. Compose this writer with [Variant] Consolidate Spark-compatible missing-value validation and errors #5977's missing-state checks and [Variant] Preserve typed-value precedence without residual pre-rewriting #5980's precedence policy.The Spark reconstruction and builder contract was checked on 4.0.4, 4.1.3 and 4.2.0. Those profiles use UTF-16 object-key ordering and clear the metadata sorted flag. #5474 tracks changing ordering when supported Spark versions permit it.
After integration, remove
rebuild_spark_variantand its replaced local writer helpers. A final BinaryView-to-Binary layout conversion is acceptable. Retain storage normalization, empty-key repair and other preparation until their own replacements land (#5477); Arrow/DataFusion upgrades are currently deferred.Completion and performance
Keep exact value and metadata comparisons for scalar-width boundaries, 63/64-byte strings, typed/residual NaNs, signed zero, dictionary traversal and ID/offset widths, Unicode and empty keys, nested residual containers, unused keys and parent nulls. Run native scan tests and unchanged upstream Spark assertions for all supported profiles.
Compare the merged baseline, candidate and vanilla Spark on matched canonical, partially shredded, fully shredded and empty-key files. Consume both byte columns, run both reader orders, and report allocation traffic separately from scan time. The target is faster whole-value reads than vanilla Spark across these inputs; keep benchmark results in the implementing PR description.
Incremental optimizations may land first. Close this issue only when the extra reconstruction pass is removed for all currently supported inputs and the compatibility and performance targets are met.
Parent: #5477; roadmap: #5438.