Skip to content

fix(spark)!: preserve partition values in file scans - #1270

Open
bvolpato wants to merge 1 commit into
substrait-io:mainfrom
bvolpato:bvolpato/fix-spark-partition-values
Open

bvolpato wants to merge 1 commit into
substrait-io:mainfrom
bvolpato:bvolpato/fix-spark-partition-values

Conversation

@bvolpato

@bvolpato bvolpato commented Sep 3, 2026 •

Copy link
Copy Markdown
Member

A partitioned Parquet read containing rows (1,10) and (2,20) round-trips as (1,NULL) and (2,NULL). The exporter lists leaf files with the full table schema, while the importer treats virtual partition columns as physical file columns.

Export each selected partition as a scan of its data schema with typed partition literals, preserve the merged column order through emit, and combine partitions with UNION_ALL. This uses Spark's resolved values rather than reconstructing them from paths, including nulls, explicit types, overlapping columns, pruned indexes, and multiple roots.

Serialize canonical file URIs and decode them as URIs so escaped partition values address the original files. Local paths with unescaped spaces remain accepted. The representation uses standard Substrait v0.102.0 relations and adds one scan/project branch per nonempty partition.

BREAKING CHANGE: Spark 3.4 conversion now rejects partitioned file scans whose file and partition column names differ only by case under case-insensitive resolution. Resolve the schema collision before conversion; this prevents changing the source plan's values to the different partition values.

Summary by CodeRabbit

  • Improvements
    • Partitioned file reads now preserve partition values and Spark’s column ordering when converted.
    • File paths containing spaces are handled, and partitioned reads support Parquet, ORC, and CSV, including multiple roots and typed, null, or escaped partition values.
  • Compatibility
    • On Spark versions that do not support case-insensitive overlap between file and partition columns, conversion now reports an unsupported operation when such columns differ only by case.

@bvolpato
bvolpato force-pushed the bvolpato/fix-spark-partition-values branch from 560db29 to b93bd7e Compare September 4, 2026 16:32
A partitioned Parquet read containing rows (1,10) and (2,20) round-trips as (1,NULL) and (2,NULL). The exporter lists leaf files with the full table schema, while the importer treats virtual partition columns as physical file columns.

Export each selected partition as a scan of its data schema with typed partition literals, preserve the merged column order through emit, and combine partitions with UNION_ALL. This uses Spark's resolved values rather than reconstructing them from paths, including nulls, explicit types, overlapping columns, pruned indexes, and multiple roots.

Serialize canonical file URIs and decode them as URIs so escaped partition values address the original files. Local paths with unescaped spaces remain accepted. The representation uses standard Substrait v0.102.0 relations and adds one scan/project branch per nonempty partition.

BREAKING CHANGE: Spark 3.4 conversion now rejects partitioned file scans whose file and partition column names differ only by case under case-insensitive resolution. Resolve the schema collision before conversion; this prevents changing the source plan's values to the different partition values.
@bvolpato
bvolpato force-pushed the bvolpato/fix-spark-partition-values branch from b93bd7e to ec2bec4 Compare September 29, 2026 05:24
@coderabbitai

coderabbitai Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: defaults

Review profile: CHILL

Plan: Advanced

Run ID: ef0ee742-16a5-44aa-8513-fafea5beee61

📥 Commits

Reviewing files that changed from the base of the PR and between 4d21734 and ec2bec4.

📒 Files selected for processing (5)
  • spark/src/main/scala/io/substrait/spark/compat/SparkCompat.scala
  • spark/src/main/scala/io/substrait/spark/logical/ToLogicalPlan.scala
  • spark/src/main/scala/io/substrait/spark/logical/ToSubstraitRel.scala
  • spark/src/main/spark-3.4/io/substrait/spark/compat/SparkCompatImpl.scala
  • spark/src/test/scala/io/substrait/spark/PartitionedFilesSuite.scala

Included review availability: This review used your included allowance. Your plan provides up to 2 included reviews per hour; 1 remain after this review.


📝 Walkthrough

Walkthrough

Partitioned Hadoop file relations now convert to partition-specific Substrait reads with partition values and merged output columns. LocalFiles paths use URI parsing with a fallback for invalid URI syntax. New tests cover partition values, overlap behavior, selected partitions, and paths containing spaces.

Changes

Partitioned file scans

Layer / File(s) Summary
Partitioned scan conversion
spark/src/main/scala/io/substrait/spark/compat/SparkCompat.scala, spark/src/main/scala/io/substrait/spark/logical/ToSubstraitRel.scala, spark/src/main/spark-3.4/io/substrait/spark/compat/SparkCompatImpl.scala, spark/src/test/scala/io/substrait/spark/PartitionedFilesSuite.scala
Partitioned Hadoop file relations produce per-partition LocalFiles reads with partition literals and merged-column remapping. Empty partition listings produce an empty scan; multiple partitions are combined with UNION_ALL. Compatibility settings and tests cover case-only overlaps, partition values, selected partitions, and empty listings.
LocalFiles URI path reconstruction
spark/src/main/scala/io/substrait/spark/logical/ToLogicalPlan.scala, spark/src/test/scala/io/substrait/spark/PartitionedFilesSuite.scala
LocalFiles paths are parsed as URIs before constructing Hadoop paths. URI syntax errors fall back to the original path string. Tests cover paths containing spaces.

Priority: ➖ Normal

Estimated code review effort: 3 (Moderate) | ~25 minutes

Change: Bug fix

Sequence Diagram(s)

sequenceDiagram
  participant SparkPlan
  participant ToSubstraitRel
  participant HadoopFsRelation
  participant SubstraitRelation
  SparkPlan->>ToSubstraitRel: convert HadoopFsRelation
  ToSubstraitRel->>HadoopFsRelation: obtain selected partitions
  HadoopFsRelation-->>ToSubstraitRel: return partition listings
  loop For each partition
    ToSubstraitRel->>SubstraitRelation: create LocalFiles read with partition literals and remap
  end
  ToSubstraitRel->>SubstraitRelation: produce empty scan, one relation, or UNION_ALL
Loading

Suggested reviewers: nielspardon

Merge Risk: 🟡 Moderate · up to ec2be

Partition-heavy reads now create separate scans for every partition, while some previously accepted raw file paths can resolve to the wrong location. Resolve or explicitly accept these scaling and compatibility regressions before merging.

Security Architecture Review

Security architecture risk: 🔵 Low · up to ec2be

The change preserves partition values while retaining the existing Spark filesystem access path. No new privilege expansion is demonstrated. Remaining uncertainty concerns legacy path compatibility, partition-driven plan growth, and the trust and filesystem credentials of production callers.

Retained concerns
No architecture-level concerns identified.

Security review details

Security Blast Radius

  • inferred — If a caller can control LocalFiles paths, those paths reach Spark filesystem indexing through the supplied Spark session. Potential exposure depends on configured filesystems, credentials, and caller validation. Maximum tenant, data-store, and environment scope cannot be determined from the available evidence; the same code-level sink existed before this PR.

Trust Boundaries and Controls

  • observed — Base and head both pass Hadoop Path objects into the same Spark file-index construction boundary. Neither inspected path applies an independent scheme or authority allowlist. This pre-existing delegation is not evidence that the PR introduced arbitrary remote access or increased privileges.
  • inferred — The fallback does not cover legacy raw strings that are syntactically valid URIs, including percent escapes, query syntax, or fragment syntax. Their former file identity is not guaranteed by the new parsing branch. The documented contract expects URIs, and no affected production input or resulting unauthorized access is established.
🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Title check ✅ Passed The title is concise, specific, and clearly describes the main change: preserving partition values in Spark file scans. It also uses the required Conventional Commit format with a breaking-change mark…
Description check ✅ Passed The description explains the defect, implementation approach, supported cases, serialization behavior, standard version, and breaking-change impact. It provides the rationale required by the repositor…
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check. Docstring coverage is scoped to functions touched by this diff. Analyzed 0 functions across 0…
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create a new PR

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@bvolpato
bvolpato marked this pull request as ready for review September 29, 2026 05:32
@nielspardon

This comment was marked as resolved.

@nielspardon nielspardon left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks — encoding the partition values as literals is the only way to express this with v0.102.0 LocalFiles until the spec closes the gap (substrait-io/substrait#1171). The main question is how the encoding scales with the number of partitions (comment on line 597); the other comments are smaller correctness fixes.

.build()
case 1 => partitions.head
case _ =>
relation.Set.builder().setOp(SetOp.UNION_ALL).addAllInputs(partitions.toSeq.asJava).build()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please collapse this back into a single scan on import: in ToLogicalPlan, turn a UNION_ALL of literal-only Projects over same-schema LocalFiles into one HadoopFsRelation, with a FileIndex that returns one PartitionDirectory per branch (like the test's withPartitions helper). As it stands each partition becomes its own Spark scan: 1000 single-file partitions plan 1000 scans and 1000 tasks instead of 32, collect takes 14 s instead of 0.2 s, and at a 512 MB driver heap it runs out of memory. The exported plan also grows with partitions × columns (about 4.5 MB for 1000 × 200, over gRPC's 4 MiB default), and listFiles(Nil, Nil) on line 553 exports every partition even under filter("part = 7") — passing partition-only filter conjuncts through to listFiles would fix that.

Comment on lines +466 to +469
new Path(new URI(path))
} catch {
// Preserve support for unescaped local paths, such as filenames containing spaces.
case _: URISyntaxException => new Path(path)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Only treat the string as a URI when it has a scheme, and rebuild the Path from its components (new Path(u.getScheme, u.getAuthority, u.getPath)). As written, a raw path that happens to parse as a URI gets reinterpreted: /data/k=a%3Ab/… decodes to a file that doesn't exist, # and ? cut the path short, and "" reads the working directory, all silently where main read the right files. A folder URI with a trailing slash (file:///tmp/t/, which File.toURI produces) also reads zero rows now.

val outputMapping = fsRelation.schema.map {
field =>
val partitionIndex =
fsRelation.partitionSchema.indexWhere(p => resolver(p.name, field.name))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Derive this mapping by position from fsRelation.schema instead of re-matching names with the resolver. Spark merged that schema when the relation was built, using toLowerCase(Locale.ROOT) and the conf at that time, so re-matching can disagree: a file column İl under il=34 gets the partition value on 3.5/4.0 although Spark keeps them separate, and changing spark.sql.caseSensitive between building and converting routes literals into the wrong column. If a resolver is still needed, use fsRelation.sparkSession.sessionState.conf.resolver — .toBoolean on line 522 throws on a padded value like "false " that Spark accepts.

val values = fsRelation.partitionSchema.zipWithIndex.map {
case (field, index) =>
ToSubstraitLiteral(
Literal(partition.values.get(index, field.dataType), field.dataType),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fit decimal values to the field's type before building the literal (Decimal.toPrecision(dt.precision, dt.scale), and a typed null when it returns null). Spark keeps the directory string's own scale here and rounds only when it writes the row, so part=1.25 read as DECIMAL(10,1) comes back as 1.2 instead of 1.3, and part=123.4 as DECIMAL(3,1) fails at collect() instead of returning null.

Comment on lines +529 to +535
if (
!SparkCompat.instance.supportsCaseInsensitivePartitionOverlap &&
fsRelation.dataSchema.exists(
data =>
fsRelation.partitionSchema.exists(
partition => data.name != partition.name && resolver(data.name, partition.name)))
) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could this reject only the queries that actually differ? On 3.4 Spark answers select("id"), select(id, value) and even the full read deterministically, and main round-tripped those to the same rows; only a filter on the partition column diverges (Spark prunes by the partition value but outputs the file value).

Literal(partition.values.get(index, field.dataType), field.dataType),
Some(field.nullable))
}
relation.Project

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Carry the partition column names somewhere in the rel tree. Nothing below the root names them, so new ToLogicalPlan(spark).convert(new ToSubstraitRel().visit(plan)) names them after the first branch's literal ([id, value, 20]), and select("part") fails; only convert(Plan) repairs this from the root names. A hint alone isn't enough, because ToLogicalPlan.visit(Project) applies hint names only when their count equals the number of expressions.

.builder()
.initialSchema(dataSchema)
.addAllItems(
partition.files

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Share one FileOrFiles builder with buildLocalFileScan (treating a flat table as a single partition). The copies already differ: the flat path lists inputFiles and gives every item fsRelation.sizeInBytes as its length, while this one uses listFiles and file.getLen.


val converted = new ToLogicalPlan(spark).convert(decoded)
assert(
DataType.equalsStructurallyByName(plan.schema, converted.schema, caseSensitiveResolution))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Compare (name, dataType) pairs here and add a PartitionDirectory(values, Nil) case. equalsStructurallyByName ignores leaf types and Row equality ignores numeric widths, so emitting a LongType partition value as an i32 literal leaves every test green, and so does deleting .filter(_.files.nonEmpty) on line 553.

@coderabbitai

This comment was marked as resolved.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants