Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 17 additions & 35 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

19 changes: 19 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,25 @@ rust-version = "1.94.0"
# Define DataFusion version
version = "55.1.0"

[patch.crates-io]
arrow = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-arith = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-array = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-avro = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-buffer = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-cast = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-cmp = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-csv = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-data = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-flight = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-ipc = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-json = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-ord = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-row = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-schema = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-select = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }
arrow-string = { git = "https://github.com/Rachelint/arrow-rs.git", rev = "cd9284f884d05e329acb293e1e95353ce555315e" }

[workspace.dependencies]
# We turn off default-features for some dependencies here so the workspaces which inherit them can
# selectively turn them on if needed, since we can override default-features = true (from false)
Expand Down
16 changes: 16 additions & 0 deletions datafusion/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1125,6 +1125,22 @@ config_namespace! {
/// aggregation ratio check and trying to switch to skipping aggregation mode
pub skip_partial_aggregation_probe_rows_threshold: usize, default = 100_000

/// (experimental) Number of groups above which a hash aggregation
/// stops growing a single hash table, so that its tables stay small
/// enough to be cache friendly. A partial aggregation then emits the
/// state of its table and starts over, as long as the emitted groups
/// do not come back. A final aggregation splits the groups seen so
/// far and all further input into hash buckets, which are aggregated
/// one after another and can be spilled and released independently;
/// it does so at a quarter of this number when its input holds about
/// one row per group. Aggregations of millions of groups per partition
/// run faster and with less memory. Moving rows into buckets has a
/// cost of its own, so aggregations that end at a few times this
/// number of groups, string keys in particular, can run a few percent
/// slower, and input that repeats its groups can use more memory.
/// Set to 0 to disable.
pub hash_aggregate_bucket_threshold: usize, default = 0

/// Should DataFusion use row number estimates at the input to decide
/// whether increasing parallelism is beneficial or not. By default,
/// only exact row numbers (not estimates) are used for this decision.
Expand Down
62 changes: 61 additions & 1 deletion datafusion/physical-expr-common/src/binary_view_map.rs
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,8 @@ where
views: Vec<u128>,
/// In-progress buffer for out-of-line string data
in_progress: Vec<u8>,
/// Maximum payload block size; extended after a reusable Partial flush.
payload_block_limit: usize,
/// Completed buffers containing string data
completed: Vec<Buffer>,

Expand Down Expand Up @@ -190,6 +192,7 @@ where
initial_map_capacity: map_capacity,
views: Vec::new(),
in_progress: Vec::new(),
payload_block_limit: BYTE_VIEW_MAX_BLOCK_SIZE,
completed: Vec::new(),
random_state: RandomState::default(),
hashes_buffer: vec![],
Expand All @@ -206,6 +209,29 @@ where
new_self
}

/// Emits all keys while retaining the hash allocation for a Partial flush.
/// The returned array owns the emitted strings. The empty map uses one payload
/// buffer, reserved from the previous payload length, up to the u32 offset
/// limit. Normal output and memory-pressure emission should use [`Self::take`].
pub fn take_state_reusing_allocation(&mut self) -> ArrayRef {
let payload_bytes = self
.completed
.iter()
.fold(self.in_progress.len(), |total, buffer| {
total.saturating_add(buffer.len())
});
let mut outgoing = Self::with_capacity(self.output_type, 0);
std::mem::swap(self, &mut outgoing);
std::mem::swap(&mut self.map, &mut outgoing.map);
self.map.clear();
self.initial_map_capacity = outgoing.initial_map_capacity;
std::mem::swap(&mut self.random_state, &mut outgoing.random_state);
self.payload_block_limit = u32::MAX as usize;
self.in_progress
.reserve_exact(payload_bytes.min(self.payload_block_limit));
outgoing.into_state()
}

/// Empties this map and releases every allocation it holds, so
/// [`Self::size`] drops to approximately zero.
///
Expand Down Expand Up @@ -538,7 +564,7 @@ where
make_view(value, 0, 0)
} else {
// Ensure buffer is big enough
if self.in_progress.len() + len > BYTE_VIEW_MAX_BLOCK_SIZE {
if self.in_progress.len() + len > self.payload_block_limit {
let flushed = std::mem::replace(
&mut self.in_progress,
Vec::with_capacity(BYTE_VIEW_MAX_BLOCK_SIZE),
Expand Down Expand Up @@ -890,6 +916,40 @@ mod tests {
assert_eq!(lazy.map.capacity(), 0);
}

#[test]
fn partial_flush_reuses_table_and_preserves_emitted_strings() {
let mut map = ArrowBytesViewMap::<usize>::new(OutputType::Utf8View);
let long_value = "x".repeat(BYTE_VIEW_MAX_BLOCK_SIZE + 1);
let expected = vec![Some("inline"), None, Some(long_value.as_str())];
let values: ArrayRef = Arc::new(StringViewArray::from(expected.clone()));
let mut emitted = Vec::new();
for _ in 0..3 {
let mut next_group = 0;
let mut groups = Vec::new();
map.insert_if_new(
&values,
|_| {
let group = next_group;
next_group += 1;
group
},
|group| groups.push(group),
);
assert_eq!(groups, vec![0, 1, 2]);
let capacity = map.map.capacity();
emitted.push(map.take_state_reusing_allocation());
assert!(map.is_empty());
assert_eq!(map.map.capacity(), capacity);
assert!(map.in_progress.capacity() >= long_value.len());
}
for array in emitted {
assert_eq!(array.as_string_view().iter().collect::<Vec<_>>(), expected);
}
map.clear_and_release();
assert_eq!(map.map.allocation_size(), 0);
assert_eq!(map.in_progress.capacity(), 0);
}

#[test]
fn clear_and_release_frees_the_preallocation_that_take_keeps() {
let mut map = ArrowBytesViewMap::<()>::with_capacity(
Expand Down
Loading