Skip to content

[Variant] Remove the extra Spark byte-reconstruction pass #5978

Description

@peterxcli

#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.

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. Preserve original value/metadata bytes only when the schema has no 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_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.

Parent: #5477; roadmap: #5438.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    area:scanParquet scan / data readingenhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions