jayzhan211 commented on code in PR #24522:
URL: https://github.com/apache/datafusion/pull/24522#discussion_r3835076392
##########
datafusion/datasource-parquet/src/projection_read_plan.rs:
##########
@@ -621,130 +628,91 @@ fn build_read_plan_with_cast_clipping(
// arrow schema). If not, never risk a wrong mask: read the whole
// root.
if root_leaves.len() != count_leaves(physical_type) {
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
continue;
}
match clip_for_cast(physical_type, &access.target_type) {
Some((kept_offsets, _pruned_type)) => {
- kept_offsets_by_root
+ if let RootRead::Partial(offsets) = root_reads
.entry(root)
- .or_default()
- .extend(kept_offsets);
+ .or_insert_with(|| RootRead::Partial(BTreeSet::new()))
+ {
+ offsets.extend(kept_offsets);
+ }
}
// Nothing prunable for this cast: every leaf is consumed.
None => {
- kept_offsets_by_root.remove(&root);
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
}
}
}
- // Add leaves reached through `get_field` to cast roots. The resolver
- // returns absolute Parquet leaf indices; convert them back to offsets in
- // their root so they can share the same union as cast clipping.
+ // Add every `get_field` root before resolving leaves. If an access matches
+ // no leaf, finalization safely falls back to a full read for that root.
+ for access in struct_accesses {
+ root_reads
+ .entry(access.root_index)
+ .or_insert_with(|| RootRead::Partial(BTreeSet::new()));
+ }
+
+ // The resolver returns absolute Parquet leaf indices. Convert each
selected
+ // leaf to a root-relative offset so casts and field accesses share one
union.
let struct_access_tree = StructAccessTree::from_accesses(struct_accesses);
for leaf in resolve_struct_field_leaves(&struct_access_tree, schema_descr)
{
let root = schema_descr.get_column_root_idx(leaf);
- if !kept_offsets_by_root.contains_key(&root) {
+ let Some(RootRead::Partial(offsets)) = root_reads.get_mut(&root) else {
continue;
- }
+ };
let Some(offset) = leaves_by_root
.get(&root)
.and_then(|root_leaves| root_leaves.binary_search(&leaf).ok())
else {
- kept_offsets_by_root.remove(&root);
- fallback_roots.insert(root);
- continue;
- };
- kept_offsets_by_root
- .get_mut(&root)
- .expect("root presence checked above")
- .insert(offset);
- }
-
- // Derive the reader's one emitted Arrow type from each merged leaf set.
- // Any unsupported partial wrapper retains the total fallback guarantee.
- let mut clipped_by_root: BTreeMap<usize, (Vec<usize>, DataType)> =
BTreeMap::new();
- for (root, kept_offsets) in kept_offsets_by_root {
- if fallback_roots.contains(&root) {
- continue;
- }
- let physical_type = file_schema.field(root).data_type();
- let root_leaves = leaves_by_root.get(&root).map_or(&[][..],
Vec::as_slice);
- let kept_offsets = kept_offsets.into_iter().collect::<Vec<_>>();
- let Some(pruned_type) = type_for_leaf_subset(physical_type,
&kept_offsets) else {
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
continue;
};
- let absolute = kept_offsets
- .into_iter()
- .map(|offset| root_leaves[offset])
- .collect();
- clipped_by_root.insert(root, (absolute, pruned_type));
+ offsets.insert(offset);
}
- // `get_field` accesses on roots not already read in full (as a whole
- // column, or as a cast that fell back) keep the existing (non-cast) leaf
- // resolution.
- let get_field_accesses: Vec<StructFieldAccess> = struct_accesses
- .iter()
- .filter(|a| {
- !whole_roots.contains(&a.root_index)
- && !fallback_roots.contains(&a.root_index)
- && !clipped_by_root.contains_key(&a.root_index)
- })
- .cloned()
- .collect();
-
let mut leaf_indices: Vec<usize> = Vec::new();
- let mut fields: BTreeMap<usize, Arc<Field>> = BTreeMap::new();
-
- for root in whole_roots.iter().chain(fallback_roots.iter()) {
- // A root with no parquet leaves contributes nothing to the mask;
- // `ProjectionMask::roots` handles that case the same way, so match it
- // rather than indexing and panicking.
- if let Some(leaves) = leaves_by_root.get(root) {
- leaf_indices.extend(leaves.iter().copied());
+ let mut fields = Vec::with_capacity(root_reads.len());
+ for (root, read) in root_reads {
+ let field = file_schema.field(root);
+ let root_leaves = leaves_by_root.get(&root).map_or(&[][..],
Vec::as_slice);
+ match read {
+ RootRead::Partial(offsets)
+ if root_leaves.len() == count_leaves(field.data_type()) =>
+ {
+ let offsets = offsets.into_iter().collect::<Vec<_>>();
+ if let Some(projected_type) =
+ type_for_leaf_subset(field.data_type(), &offsets)
+ {
+ leaf_indices
+ .extend(offsets.into_iter().map(|offset|
root_leaves[offset]));
+ fields.push(field_with_type(field, projected_type));
Review Comment:
`assemble_read_plan` updated
##########
datafusion/datasource-parquet/src/projection_read_plan.rs:
##########
@@ -621,130 +628,91 @@ fn build_read_plan_with_cast_clipping(
// arrow schema). If not, never risk a wrong mask: read the whole
// root.
if root_leaves.len() != count_leaves(physical_type) {
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
continue;
}
match clip_for_cast(physical_type, &access.target_type) {
Some((kept_offsets, _pruned_type)) => {
- kept_offsets_by_root
+ if let RootRead::Partial(offsets) = root_reads
.entry(root)
- .or_default()
- .extend(kept_offsets);
+ .or_insert_with(|| RootRead::Partial(BTreeSet::new()))
+ {
+ offsets.extend(kept_offsets);
+ }
}
// Nothing prunable for this cast: every leaf is consumed.
None => {
- kept_offsets_by_root.remove(&root);
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
}
}
}
- // Add leaves reached through `get_field` to cast roots. The resolver
- // returns absolute Parquet leaf indices; convert them back to offsets in
- // their root so they can share the same union as cast clipping.
+ // Add every `get_field` root before resolving leaves. If an access matches
+ // no leaf, finalization safely falls back to a full read for that root.
+ for access in struct_accesses {
+ root_reads
+ .entry(access.root_index)
+ .or_insert_with(|| RootRead::Partial(BTreeSet::new()));
+ }
+
+ // The resolver returns absolute Parquet leaf indices. Convert each
selected
+ // leaf to a root-relative offset so casts and field accesses share one
union.
let struct_access_tree = StructAccessTree::from_accesses(struct_accesses);
for leaf in resolve_struct_field_leaves(&struct_access_tree, schema_descr)
{
let root = schema_descr.get_column_root_idx(leaf);
- if !kept_offsets_by_root.contains_key(&root) {
+ let Some(RootRead::Partial(offsets)) = root_reads.get_mut(&root) else {
continue;
- }
+ };
let Some(offset) = leaves_by_root
.get(&root)
.and_then(|root_leaves| root_leaves.binary_search(&leaf).ok())
else {
- kept_offsets_by_root.remove(&root);
- fallback_roots.insert(root);
- continue;
- };
- kept_offsets_by_root
- .get_mut(&root)
- .expect("root presence checked above")
- .insert(offset);
- }
-
- // Derive the reader's one emitted Arrow type from each merged leaf set.
- // Any unsupported partial wrapper retains the total fallback guarantee.
- let mut clipped_by_root: BTreeMap<usize, (Vec<usize>, DataType)> =
BTreeMap::new();
- for (root, kept_offsets) in kept_offsets_by_root {
- if fallback_roots.contains(&root) {
- continue;
- }
- let physical_type = file_schema.field(root).data_type();
- let root_leaves = leaves_by_root.get(&root).map_or(&[][..],
Vec::as_slice);
- let kept_offsets = kept_offsets.into_iter().collect::<Vec<_>>();
- let Some(pruned_type) = type_for_leaf_subset(physical_type,
&kept_offsets) else {
- fallback_roots.insert(root);
+ root_reads.insert(root, RootRead::Full);
continue;
};
- let absolute = kept_offsets
- .into_iter()
- .map(|offset| root_leaves[offset])
- .collect();
- clipped_by_root.insert(root, (absolute, pruned_type));
+ offsets.insert(offset);
}
- // `get_field` accesses on roots not already read in full (as a whole
- // column, or as a cast that fell back) keep the existing (non-cast) leaf
- // resolution.
- let get_field_accesses: Vec<StructFieldAccess> = struct_accesses
- .iter()
- .filter(|a| {
- !whole_roots.contains(&a.root_index)
- && !fallback_roots.contains(&a.root_index)
- && !clipped_by_root.contains_key(&a.root_index)
- })
- .cloned()
- .collect();
-
let mut leaf_indices: Vec<usize> = Vec::new();
- let mut fields: BTreeMap<usize, Arc<Field>> = BTreeMap::new();
-
- for root in whole_roots.iter().chain(fallback_roots.iter()) {
- // A root with no parquet leaves contributes nothing to the mask;
- // `ProjectionMask::roots` handles that case the same way, so match it
- // rather than indexing and panicking.
- if let Some(leaves) = leaves_by_root.get(root) {
- leaf_indices.extend(leaves.iter().copied());
+ let mut fields = Vec::with_capacity(root_reads.len());
+ for (root, read) in root_reads {
+ let field = file_schema.field(root);
+ let root_leaves = leaves_by_root.get(&root).map_or(&[][..],
Vec::as_slice);
+ match read {
+ RootRead::Partial(offsets)
+ if root_leaves.len() == count_leaves(field.data_type()) =>
Review Comment:
remove the previous check
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]