Conversation
560db29 to
b93bd7e
Compare
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.
b93bd7e to
ec2bec4
Compare
|
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 configurationConfiguration used: defaults Review profile: CHILL Plan: Advanced Run ID: 📒 Files selected for processing (5)
Included review availability: This review used your included allowance. Your plan provides up to 2 included reviews per hour; 1 remain after this review. 📝 WalkthroughWalkthroughPartitioned 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. ChangesPartitioned file scans
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
Suggested reviewers: Merge Risk: 🟡 Moderate · up to 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 ReviewSecurity architecture risk: 🔵 Low · up to 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 Security review detailsSecurity Blast Radius
Trust Boundaries and Controls
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
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. Comment |
This comment was marked as resolved.
This comment was marked as resolved.
nielspardon
left a comment
There was a problem hiding this comment.
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() |
There was a problem hiding this comment.
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.
| new Path(new URI(path)) | ||
| } catch { | ||
| // Preserve support for unescaped local paths, such as filenames containing spaces. | ||
| case _: URISyntaxException => new Path(path) |
There was a problem hiding this comment.
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)) |
There was a problem hiding this comment.
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), |
There was a problem hiding this comment.
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.
| if ( | ||
| !SparkCompat.instance.supportsCaseInsensitivePartitionOverlap && | ||
| fsRelation.dataSchema.exists( | ||
| data => | ||
| fsRelation.partitionSchema.exists( | ||
| partition => data.name != partition.name && resolver(data.name, partition.name))) | ||
| ) { |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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)) |
There was a problem hiding this comment.
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.
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