Skip to content

Commit cf6819d

Browse files
Rachel IntRachel Int
authored andcommitted
perf: retain gather when scattering StringView bucket ranges
1 parent 085c221 commit cf6819d

3 files changed

Lines changed: 45 additions & 64 deletions

File tree

‎Cargo.lock‎

Lines changed: 17 additions & 17 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 17 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -84,23 +84,23 @@ rust-version = "1.94.0"
8484
version = "55.1.0"
8585

8686
[patch.crates-io]
87-
arrow = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
88-
arrow-arith = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
89-
arrow-array = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
90-
arrow-avro = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
91-
arrow-buffer = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
92-
arrow-cast = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
93-
arrow-cmp = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
94-
arrow-csv = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
95-
arrow-data = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
96-
arrow-flight = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
97-
arrow-ipc = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
98-
arrow-json = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
99-
arrow-ord = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
100-
arrow-row = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
101-
arrow-schema = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
102-
arrow-select = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
103-
arrow-string = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "8d0558fc6f6a2acb47e91207ef7b4acf24a72245" }
87+
arrow = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
88+
arrow-arith = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
89+
arrow-array = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
90+
arrow-avro = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
91+
arrow-buffer = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
92+
arrow-cast = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
93+
arrow-cmp = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
94+
arrow-csv = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
95+
arrow-data = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
96+
arrow-flight = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
97+
arrow-ipc = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
98+
arrow-json = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
99+
arrow-ord = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
100+
arrow-row = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
101+
arrow-schema = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
102+
arrow-select = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
103+
arrow-string = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
104104

105105
[workspace.dependencies]
106106
# We turn off default-features for some dependencies here so the workspaces which inherit them can

‎datafusion/physical-plan/src/aggregates/final_buckets.rs‎

Lines changed: 11 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -150,8 +150,8 @@ impl ColumnArray {
150150
}
151151
}
152152

153-
/// Builds full batches by scattering StringView columns and coalescing the
154-
/// remaining columns after gathering them into bucket order. Source ranges use
153+
/// Builds full batches by gathering columns into bucket order, then scattering
154+
/// their contiguous ranges into in-progress target arrays. Source ranges use
155155
/// the same copy and compaction decisions as `BatchCoalescer`, without
156156
/// allocating a small record batch for every destination.
157157
struct ColumnBatchBuilder {
@@ -412,39 +412,16 @@ impl FinalBuckets {
412412
*cursor += 1;
413413
}
414414

415-
// Scatter StringView columns directly from source order into every
416-
// destination. Other columns are gathered once in bucket order.
415+
// Gather every column in bucket order, then scatter each StringView
416+
// column's contiguous ranges into all in-memory destinations at once.
417417
// Spilled buckets keep the existing gather-and-slice path.
418418
let indices: PrimitiveArray<UInt32Type> =
419419
std::mem::take(&mut self.reordered_indices).into();
420+
let columns = take_arrays(batch.columns(), &indices, None)?;
420421
let can_scatter = self
421422
.buckets
422423
.iter()
423-
.all(|bucket| bucket.spill_file.is_none())
424-
&& batch
425-
.columns()
426-
.iter()
427-
.any(|source| source.data_type() == &DataType::Utf8View);
428-
let columns = if can_scatter {
429-
batch
430-
.columns()
431-
.iter()
432-
.map(|source| {
433-
if source.data_type() == &DataType::Utf8View {
434-
Ok(Arc::clone(source))
435-
} else {
436-
compute::take(source.as_ref(), &indices, None)
437-
}
438-
})
439-
.collect::<ArrowResult<Vec<_>>>()?
440-
} else {
441-
take_arrays(batch.columns(), &indices, None)?
442-
};
443-
let bucket_ids: Vec<_> = if can_scatter {
444-
self.hashes.iter().map(|&hash| bucket_of(hash)).collect()
445-
} else {
446-
Vec::new()
447-
};
424+
.all(|bucket| bucket.spill_file.is_none());
448425
for (column_index, source) in columns.iter().enumerate() {
449426
if can_scatter && source.data_type() == &DataType::Utf8View {
450427
let mut targets: Vec<_> = self
@@ -464,7 +441,11 @@ impl FinalBuckets {
464441
}
465442
})
466443
.collect();
467-
compute::scatter_string_views(source, &bucket_ids, &mut targets)?;
444+
compute::scatter_string_view_ranges(
445+
source,
446+
&self.bucket_sizes,
447+
&mut targets,
448+
)?;
468449
continue;
469450
}
470451
for (bucket, (&start, &size)) in self

0 commit comments

Comments
 (0)