This is an automated email from the ASF dual-hosted git repository.
tisonkun pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datasketches-rust.git
The following commit(s) were added to refs/heads/main by this push:
new c3c08f9 perf(tdigest): optimize partial aggregate compute path (#231)
c3c08f9 is described below
commit c3c08f9b5ce0a537b5a1c77b4ea11db8e3b5f2cb
Author: tison <[email protected]>
AuthorDate: Thu Aug 27 17:31:05 2026 +0800
perf(tdigest): optimize partial aggregate compute path (#231)
---
CHANGELOG.md | 2 +-
benchmarks/tdigest/merge.rs | 37 ++
datasketches/src/tdigest/sketch.rs | 462 ++++++++++---------------
tests-integration/tests/serde_tests/tdigest.rs | 21 +-
tests-integration/tests/tdigest_test/sketch.rs | 6 +-
5 files changed, 241 insertions(+), 287 deletions(-)
diff --git a/CHANGELOG.md b/CHANGELOG.md
index eba1dae..6b4ffd6 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -14,7 +14,7 @@ All significant changes to this project will be documented in
this file.
### Performance improvements
-* Reduce T-Digest allocation overhead and retained memory across updates,
compression, merges, serialization, deserialization, and freezing while
preserving the serialized format.
+* Reduce T-Digest allocation overhead and retained memory across updates,
compression, merges, serialization, deserialization, and freezing; linearly
merge sorted centroid buffers and decode validated native payloads directly
while preserving the serialized format.
* Reduce CPC serialization and deserialization allocations by encoding
directly into the output buffer and decoding directly from the input payload.
### Bug fixes
diff --git a/benchmarks/tdigest/merge.rs b/benchmarks/tdigest/merge.rs
index b66b5cc..99de685 100644
--- a/benchmarks/tdigest/merge.rs
+++ b/benchmarks/tdigest/merge.rs
@@ -26,6 +26,7 @@ use super::support::SMALL_ROWS_PER_PARTIAL;
use super::support::build_mut_digest;
use super::support::partial_digests;
use super::support::partial_digests_with;
+use super::support::serialized_partial_digests;
use super::support::values;
#[divan::bench]
@@ -100,3 +101,39 @@ fn partials(bencher: Bencher) {
black_box(merged)
});
}
+
+#[divan::bench(args = [SMALL_ROWS_PER_PARTIAL, ROWS_PER_PARTIAL])]
+fn serialized_partials(bencher: Bencher, rows_per_partial: usize) {
+ let partials = serialized_partial_digests(64, rows_per_partial);
+
+ bencher
+ .counter(ItemsCount::new(64 * rows_per_partial))
+ .bench_local(|| {
+ let mut merged = TDigestMut::default();
+ for partial in black_box(&partials) {
+ let partial = TDigestMut::deserialize(partial, false).unwrap();
+ merged.merge(&partial);
+ }
+ black_box(merged.quantile(0.5))
+ });
+}
+
+#[divan::bench]
+fn serialized_overlapping_partials(bencher: Bencher) {
+ let values = values(64 * ROWS_PER_PARTIAL);
+ let partials = values
+ .chunks_exact(ROWS_PER_PARTIAL)
+ .map(|values| build_mut_digest(values).serialize())
+ .collect::<Vec<_>>();
+
+ bencher
+ .counter(ItemsCount::new(values.len()))
+ .bench_local(|| {
+ let mut merged = TDigestMut::default();
+ for partial in black_box(&partials) {
+ let partial = TDigestMut::deserialize(partial, false).unwrap();
+ merged.merge(&partial);
+ }
+ black_box(merged.quantile(0.5))
+ });
+}
diff --git a/datasketches/src/tdigest/sketch.rs
b/datasketches/src/tdigest/sketch.rs
index e24a99f..9acb558 100644
--- a/datasketches/src/tdigest/sketch.rs
+++ b/datasketches/src/tdigest/sketch.rs
@@ -44,83 +44,122 @@ const INITIAL_UNMERGED_CAPACITY: usize = 8;
/// Default weight for single values.
const DEFAULT_WEIGHT: NonZeroU64 = NonZeroU64::new(1).unwrap();
-// The update buffer has two physical representations:
-//
-// * `Staging` stores raw `f64` values compactly before the first compression.
-// * `Centroids` stores `[compressed prefix | unmerged unit-weight tail]` in
one allocation. The
-// tail length identifies the boundary between the two regions.
-//
-// Compression permanently transitions a non-empty buffer from `Staging` to
`Centroids`.
-#[derive(Debug, Clone)]
-enum TDigestBuffer {
- Staging(Vec<f64>),
- Centroids {
- centroids: Vec<Centroid>,
- unmerged_tail_len: usize,
- },
+// Centroids are stored as `[compressed prefix | unmerged unit-weight tail]`.
Carrying a unit weight
+// for each unmerged value lets updates, compression, and merge reuse one
allocation instead of
+// converting raw values into a second vector during compression.
+#[derive(Debug, Clone, Default)]
+struct TDigestBuffer {
+ centroids: Vec<Centroid>,
+ unmerged_tail_len: usize,
}
-impl Default for TDigestBuffer {
- fn default() -> Self {
- TDigestBuffer::Staging(vec![])
+impl TDigestBuffer {
+ fn new(centroids: Vec<Centroid>, unmerged_tail_len: usize) -> Self {
+ debug_assert!(unmerged_tail_len <= centroids.len());
+ TDigestBuffer {
+ centroids,
+ unmerged_tail_len,
+ }
}
-}
-impl TDigestBuffer {
fn len(&self) -> usize {
- match self {
- TDigestBuffer::Staging(values) => values.len(),
- TDigestBuffer::Centroids { centroids, .. } => centroids.len(),
- }
+ self.centroids.len()
}
fn is_empty(&self) -> bool {
- self.len() == 0
+ self.centroids.is_empty()
}
fn unmerged_len(&self) -> usize {
- match self {
- TDigestBuffer::Staging(values) => values.len(),
- TDigestBuffer::Centroids {
- unmerged_tail_len, ..
- } => *unmerged_tail_len,
+ self.unmerged_tail_len
+ }
+
+ fn compressed_prefix_len(&self) -> usize {
+ self.centroids.len() - self.unmerged_tail_len
+ }
+
+ fn push_unmerged(&mut self, value: f64, max_unmerged: usize) {
+ debug_assert!(self.unmerged_tail_len < max_unmerged);
+ if self.centroids.len() == self.centroids.capacity() {
+ let target_unmerged = if self.unmerged_tail_len == 0 {
+ INITIAL_UNMERGED_CAPACITY
+ } else if self.unmerged_tail_len == INITIAL_UNMERGED_CAPACITY {
+ // Once a digest outgrows a tiny group, skip an extra
allocator round trip while
+ // keeping the first allocation small.
+ (INITIAL_UNMERGED_CAPACITY * UNMERGED_MULTIPLIER *
UNMERGED_MULTIPLIER)
+ .min(max_unmerged)
+ } else {
+ self.unmerged_tail_len
+ .saturating_mul(UNMERGED_MULTIPLIER)
+ .min(max_unmerged)
+ };
+ let target_capacity =
self.compressed_prefix_len().saturating_add(target_unmerged);
+ self.centroids
+
.reserve_exact(target_capacity.saturating_sub(self.centroids.len()));
}
+
+ self.centroids.push(Centroid {
+ mean: value,
+ weight: DEFAULT_WEIGHT,
+ });
+ self.unmerged_tail_len += 1;
+ }
+
+ /// Returns all centroids in the tie order expected by stable compression
sorting.
+ ///
+ /// The buffer is rotated from `[compressed | unmerged]` to `[unmerged |
compressed]`, so new
+ /// values stay before existing centroids when their means are equal.
+ fn into_centroids_for_compression(mut self) -> Vec<Centroid> {
+ debug_assert_ne!(self.unmerged_tail_len, 0);
+ let compressed_prefix_len = self.compressed_prefix_len();
+ self.centroids.rotate_left(compressed_prefix_len);
+ self.centroids
+ }
+
+ /// Combines this buffer with a non-empty borrowed buffer in stable mean
order.
+ fn into_merged_centroids(mut self, other: &TDigestBuffer) -> Vec<Centroid>
{
+ debug_assert!(!other.is_empty(), "an empty right-hand buffer is a
no-op");
+ if self.unmerged_tail_len == 0
+ && other.unmerged_tail_len == 0
+ && centroids_are_sorted(&self.centroids)
+ && centroids_are_sorted(&other.centroids)
+ {
+ merge_sorted_centroids(&mut self.centroids, &other.centroids);
+ return self.centroids;
+ }
+
+ let compressed_prefix_len = self.compressed_prefix_len();
+ self.centroids.reserve(other.len());
+ let other_prefix_len = other.compressed_prefix_len();
+ self.centroids
+ .extend_from_slice(&other.centroids[other_prefix_len..]);
+ self.centroids
+ .extend_from_slice(&other.centroids[..other_prefix_len]);
+ // Preserve the stable tie order: left unmerged, right unmerged and
compressed, then the
+ // left compressed prefix.
+ self.centroids.rotate_left(compressed_prefix_len);
+ self.centroids.sort_by(centroid_cmp);
+ self.centroids
}
- /// Returns compressed centroids after the caller has processed staged
values.
fn compressed_centroids(&self) -> &[Centroid] {
- match self {
- TDigestBuffer::Staging(values) if values.is_empty() => &[],
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- } => centroids,
- _ => unreachable!(
- "t-digest buffer must be compressed before reading centroids:
{self:?}"
- ),
- }
+ assert_eq!(
+ self.unmerged_tail_len, 0,
+ "t-digest buffer must be compressed before reading centroids"
+ );
+ &self.centroids
}
fn into_compressed_centroids(self) -> Vec<Centroid> {
- match self {
- TDigestBuffer::Staging(values) if values.is_empty() => vec![],
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- } => centroids,
- _ => unreachable!(
- "t-digest buffer must be compressed before reading centroids:
{self:?}"
- ),
- }
+ assert_eq!(
+ self.unmerged_tail_len, 0,
+ "t-digest buffer must be compressed before reading centroids"
+ );
+ self.centroids
}
fn estimated_size(&self) -> usize {
- match self {
- TDigestBuffer::Staging(values) => values.capacity() *
size_of::<f64>(),
- TDigestBuffer::Centroids { centroids, .. } => {
- centroids.capacity() * size_of::<Centroid>()
- }
- }
+ self.centroids.capacity() * size_of::<Centroid>()
}
}
@@ -136,8 +175,8 @@ pub struct TDigestMut {
max: f64,
buffer: TDigestBuffer,
- // Weight represented by the compressed prefix. Staged values or the
unmerged tail contribute
- // one each and are counted separately by `TDigestBuffer::unmerged_len`.
+ // Weight represented by the compressed prefix. The unmerged tail
contributes one per
+ // centroid and is counted separately by `TDigestBuffer::unmerged_len`.
compressed_weight: u64,
}
@@ -170,7 +209,7 @@ impl TDigestMut {
false,
f64::INFINITY,
f64::NEG_INFINITY,
- TDigestBuffer::Staging(vec![]),
+ TDigestBuffer::default(),
0,
)
}
@@ -203,7 +242,7 @@ impl TDigestMut {
false,
f64::INFINITY,
f64::NEG_INFINITY,
- TDigestBuffer::Staging(vec![]),
+ TDigestBuffer::default(),
0,
))
}
@@ -218,13 +257,8 @@ impl TDigestMut {
compressed_weight: u64,
) -> Self {
assert!(k >= 10, "k must be at least 10");
- debug_assert!(match &buffer {
- TDigestBuffer::Staging(_) => compressed_weight == 0,
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len,
- } => *unmerged_tail_len <= centroids.len(),
- });
+ debug_assert!(buffer.unmerged_tail_len <= buffer.centroids.len());
+ debug_assert!(buffer.compressed_prefix_len() != 0 || compressed_weight
== 0);
TDigestMut {
k,
@@ -268,70 +302,10 @@ impl TDigestMut {
}
let max_unmerged = self.max_unmerged();
- if let TDigestBuffer::Staging(values) = &mut self.buffer {
- if values.len() < max_unmerged {
- if values.len() == values.capacity() {
- let target_capacity = if values.capacity() == 0 {
- INITIAL_UNMERGED_CAPACITY
- } else if values.capacity() == INITIAL_UNMERGED_CAPACITY {
- // Once a digest outgrows a tiny group, skip an extra
allocator round trip
- // while keeping the first allocation small.
- (INITIAL_UNMERGED_CAPACITY * UNMERGED_MULTIPLIER *
UNMERGED_MULTIPLIER)
- .min(max_unmerged)
- } else {
- values
- .capacity()
- .saturating_mul(UNMERGED_MULTIPLIER)
- .min(max_unmerged)
- };
-
values.reserve_exact(target_capacity.saturating_sub(values.len()));
- }
-
- values.push(value);
- self.min = self.min.min(value);
- self.max = self.max.max(value);
- return;
- }
- self.compress();
- }
-
- if matches!(
- &self.buffer,
- TDigestBuffer::Centroids { unmerged_tail_len, .. } if
*unmerged_tail_len >= max_unmerged
- ) {
+ if self.buffer.unmerged_len() >= max_unmerged {
self.compress();
}
-
- let TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len,
- } = &mut self.buffer
- else {
- unreachable!("a full staging buffer must become centroid-backed
after compression");
- };
- if centroids.len() == centroids.capacity() {
- let target_unmerged = if *unmerged_tail_len == 0 {
- INITIAL_UNMERGED_CAPACITY
- } else if *unmerged_tail_len == INITIAL_UNMERGED_CAPACITY {
- // Once a digest outgrows a tiny group, skip an extra
allocator round trip while
- // keeping the first allocation small.
- (INITIAL_UNMERGED_CAPACITY * UNMERGED_MULTIPLIER *
UNMERGED_MULTIPLIER)
- .min(max_unmerged)
- } else {
- unmerged_tail_len
- .saturating_mul(UNMERGED_MULTIPLIER)
- .min(max_unmerged)
- };
- let num_merged = centroids.len() - *unmerged_tail_len;
- let target_capacity = num_merged.saturating_add(target_unmerged);
-
centroids.reserve_exact(target_capacity.saturating_sub(centroids.len()));
- }
-
- centroids.push(Centroid {
- mean: value,
- weight: DEFAULT_WEIGHT,
- });
- *unmerged_tail_len += 1;
+ self.buffer.push_unmerged(value, max_unmerged);
self.min = self.min.min(value);
self.max = self.max.max(value);
}
@@ -388,47 +362,9 @@ impl TDigestMut {
return;
}
- let buffer = std::mem::take(&mut self.buffer);
- let (mut merge_buffer, existing_prefix_len, self_unmerged_weight) =
match buffer {
- TDigestBuffer::Staging(values) => {
- let self_unmerged_weight = values.len() as u64;
- let mut merge_buffer = Vec::with_capacity(values.len() +
other.buffer.len());
- merge_buffer.extend(values.into_iter().map(|mean| Centroid {
- mean,
- weight: DEFAULT_WEIGHT,
- }));
- (merge_buffer, 0, self_unmerged_weight)
- }
- TDigestBuffer::Centroids {
- mut centroids,
- unmerged_tail_len,
- } => {
- let existing_prefix_len = centroids.len() - unmerged_tail_len;
- centroids.reserve(other.buffer.len());
- (centroids, existing_prefix_len, unmerged_tail_len as u64)
- }
- };
- match &other.buffer {
- TDigestBuffer::Staging(values) => {
- merge_buffer.extend(values.iter().copied().map(|mean| Centroid
{
- mean,
- weight: DEFAULT_WEIGHT,
- }));
- }
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len,
- } => {
- let other_prefix_len = centroids.len() - unmerged_tail_len;
- merge_buffer.extend_from_slice(¢roids[other_prefix_len..]);
- merge_buffer.extend_from_slice(¢roids[..other_prefix_len]);
- }
- }
- // Preserve the original insertion order for equal means because
t-digest requires a stable
- // sort: this digest's buffered/unmerged values, the other digest's
buffered/unmerged values
- // and compressed centroids, then this digest's compressed centroids.
- merge_buffer.rotate_left(existing_prefix_len);
- self.compress_centroids(merge_buffer, self_unmerged_weight +
other.total_weight())
+ let self_unmerged_weight = self.buffer.unmerged_len() as u64;
+ let centroids = std::mem::take(&mut
self.buffer).into_merged_centroids(&other.buffer);
+ self.compress_sorted_centroids(centroids, self_unmerged_weight +
other.total_weight())
}
/// Converts this mutable t-digest into an immutable one.
@@ -739,13 +675,13 @@ impl TDigestMut {
reverse_merge,
value,
value,
- TDigestBuffer::Centroids {
- centroids: vec![Centroid {
+ TDigestBuffer::new(
+ vec![Centroid {
mean: value,
weight: DEFAULT_WEIGHT,
}],
- unmerged_tail_len: 0,
- },
+ 0,
+ ),
1,
));
}
@@ -775,61 +711,40 @@ impl TDigestMut {
} else {
(size_of::<f64>() + size_of::<u64>(), size_of::<f64>())
};
- let required_payload_bytes = num_centroids
+ let centroid_payload_bytes = num_centroids
.checked_mul(centroid_bytes)
- .and_then(|bytes| {
- num_buffered
- .checked_mul(buffered_value_bytes)
- .and_then(|buffered_bytes|
bytes.checked_add(buffered_bytes))
- })
.ok_or_else(|| Error::deserial("TDigest payload size exceeds the
supported size"))?;
- if cursor.remaining().len() < required_payload_bytes {
+ let buffered_payload_bytes = num_buffered
+ .checked_mul(buffered_value_bytes)
+ .ok_or_else(|| Error::deserial("TDigest payload size exceeds the
supported size"))?;
+ let required_payload_bytes = centroid_payload_bytes
+ .checked_add(buffered_payload_bytes)
+ .ok_or_else(|| Error::deserial("TDigest payload size exceeds the
supported size"))?;
+ let remaining = cursor.remaining();
+ if remaining.len() < required_payload_bytes {
return Err(Error::insufficient_data(format!(
"TDigest payload requires {required_payload_bytes} bytes, got
{}",
- cursor.remaining().len()
+ remaining.len()
)));
}
- if num_centroids == 0 {
- checked_weight_sum(0, num_buffered as u64)?;
- let mut initial_buffer = Vec::with_capacity(num_buffered);
- for _ in 0..num_buffered {
- let value = if is_f32 {
- cursor
- .read_f32_le()
- .map_err(insufficient_data("buffered_value"))? as f64
- } else {
- cursor
- .read_f64_le()
- .map_err(insufficient_data("buffered_value"))?
- };
- check_non_nan(value, "buffered_value mean")?;
- check_finite(value, "buffered_value mean")?;
- initial_buffer.push(value);
- }
- return Ok(TDigestMut::make(
- k,
- reverse_merge,
- min,
- max,
- TDigestBuffer::Staging(initial_buffer),
- 0,
- ));
- }
+ // Check the whole payload once so fixed-width records can be decoded
without per-field I/O.
+ let (centroid_payload, buffered_payload) =
+
remaining[..required_payload_bytes].split_at(centroid_payload_bytes);
let stored_centroids =
num_centroids.checked_add(num_buffered).ok_or_else(|| {
Error::deserial("num_centroids and num_buffered exceed the
supported size")
})?;
let mut centroids = Vec::with_capacity(stored_centroids);
let mut compressed_weight = 0u64;
- for _ in 0..num_centroids {
+ for bytes in centroid_payload.chunks_exact(centroid_bytes) {
let (mean, weight) = if is_f32 {
(
- cursor.read_f32_le().map_err(insufficient_data("mean"))?
as f64,
- cursor.read_u32_le().map_err(insufficient_data("weight"))?
as u64,
+ f32::from_le_bytes(bytes[..4].try_into().unwrap()) as f64,
+ u32::from_le_bytes(bytes[4..].try_into().unwrap()) as u64,
)
} else {
(
- cursor.read_f64_le().map_err(insufficient_data("mean"))?,
- cursor.read_u64_le().map_err(insufficient_data("weight"))?,
+ f64::from_le_bytes(bytes[..8].try_into().unwrap()),
+ u64::from_le_bytes(bytes[8..].try_into().unwrap()),
)
};
check_non_nan(mean, "centroid mean")?;
@@ -839,15 +754,11 @@ impl TDigestMut {
centroids.push(Centroid { mean, weight });
}
checked_weight_sum(compressed_weight, num_buffered as u64)?;
- for _ in 0..num_buffered {
+ for bytes in buffered_payload.chunks_exact(buffered_value_bytes) {
let value = if is_f32 {
- cursor
- .read_f32_le()
- .map_err(insufficient_data("buffered_value"))? as f64
+ f32::from_le_bytes(bytes.try_into().unwrap()) as f64
} else {
- cursor
- .read_f64_le()
- .map_err(insufficient_data("buffered_value"))?
+ f64::from_le_bytes(bytes.try_into().unwrap())
};
check_non_nan(value, "buffered_value mean")?;
check_finite(value, "buffered_value mean")?;
@@ -861,10 +772,7 @@ impl TDigestMut {
reverse_merge,
min,
max,
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: num_buffered,
- },
+ TDigestBuffer::new(centroids, num_buffered),
compressed_weight,
))
}
@@ -916,10 +824,7 @@ impl TDigestMut {
false,
min,
max,
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- },
+ TDigestBuffer::new(centroids, 0),
total_weight,
))
}
@@ -963,10 +868,7 @@ impl TDigestMut {
false,
min,
max,
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- },
+ TDigestBuffer::new(centroids, 0),
total_weight,
))
}
@@ -980,43 +882,13 @@ impl TDigestMut {
/// Processes unmerged values and merges centroids if needed.
fn compress(&mut self) {
- let buffer = std::mem::take(&mut self.buffer);
- match buffer {
- TDigestBuffer::Staging(values) if values.is_empty() => {
- self.buffer = TDigestBuffer::Staging(values);
- }
- TDigestBuffer::Staging(values) => {
- debug_assert_eq!(self.compressed_weight, 0);
- let weight = values.len() as u64;
- let mut centroids = Vec::with_capacity(values.len());
- centroids.extend(values.into_iter().map(|mean| Centroid {
- mean,
- weight: DEFAULT_WEIGHT,
- }));
- self.compress_centroids(centroids, weight);
- }
- TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- } => {
- // Preserve compact deserialized images verbatim, including
images with more
- // centroids than this implementation would normally produce.
- self.buffer = TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- };
- }
- TDigestBuffer::Centroids {
- mut centroids,
- unmerged_tail_len,
- } => {
- let compressed_prefix_len = centroids.len() -
unmerged_tail_len;
- // Preserve the original insertion order for equal means
because t-digest requires
- // a stable sort: unmerged values before existing centroids.
- centroids.rotate_left(compressed_prefix_len);
- self.compress_centroids(centroids, unmerged_tail_len as u64);
- }
+ let additional_weight = self.buffer.unmerged_len() as u64;
+ if additional_weight == 0 {
+ // Also preserves fully compressed deserialized images verbatim.
+ return;
}
+ let centroids = std::mem::take(&mut
self.buffer).into_centroids_for_compression();
+ self.compress_centroids(centroids, additional_weight);
}
/// Compresses the given centroids into this t-digest.
@@ -1028,10 +900,16 @@ impl TDigestMut {
/// stored in `self`.
/// * `additional_weight` is the total weight not yet included in
`self.compressed_weight`.
/// * Every centroid mean in `centroids` is finite.
- /// * `self.buffer` becomes centroid-backed with no unmerged values before
returning.
+ /// * `self.buffer` has no unmerged values before returning.
fn compress_centroids(&mut self, mut centroids: Vec<Centroid>,
additional_weight: u64) {
debug_assert!(!centroids.is_empty());
centroids.sort_by(centroid_cmp);
+ self.compress_sorted_centroids(centroids, additional_weight);
+ }
+
+ fn compress_sorted_centroids(&mut self, mut centroids: Vec<Centroid>,
additional_weight: u64) {
+ debug_assert!(!centroids.is_empty());
+ debug_assert!(centroids_are_sorted(¢roids));
if self.reverse_merge {
centroids.reverse();
}
@@ -1075,10 +953,7 @@ impl TDigestMut {
self.max = self.max.max(centroids[num_centroids - 1].mean);
self.reverse_merge = !self.reverse_merge;
self.reduce_retained_capacity(&mut centroids);
- self.buffer = TDigestBuffer::Centroids {
- centroids,
- unmerged_tail_len: 0,
- };
+ self.buffer = TDigestBuffer::new(centroids, 0);
}
fn reduce_retained_capacity(&self, centroids: &mut Vec<Centroid>) {
@@ -1304,10 +1179,7 @@ impl TDigest {
self.reverse_merge,
self.min,
self.max,
- TDigestBuffer::Centroids {
- centroids: self.centroids,
- unmerged_tail_len: 0,
- },
+ TDigestBuffer::new(self.centroids, 0),
self.centroids_weight,
)
}
@@ -1551,6 +1423,42 @@ fn centroid_cmp(a: &Centroid, b: &Centroid) -> Ordering {
}
}
+fn centroids_are_sorted(centroids: &[Centroid]) -> bool {
+ centroids
+ .windows(2)
+ .all(|pair| centroid_cmp(&pair[0], &pair[1]) != Ordering::Greater)
+}
+
+fn merge_sorted_centroids(left: &mut Vec<Centroid>, right: &[Centroid]) {
+ debug_assert!(!right.is_empty());
+ debug_assert!(centroids_are_sorted(left));
+ debug_assert!(centroids_are_sorted(right));
+
+ let mut left_index = left.len();
+ let mut right_index = right.len();
+ let mut output_index = left_index + right_index;
+ left.reserve(right.len());
+ left.resize(output_index, right[0]);
+
+ while left_index > 0 && right_index > 0 {
+ let left_centroid = left[left_index - 1];
+ let right_centroid = right[right_index - 1];
+ output_index -= 1;
+ // Taking the left side on ties while filling backward keeps the right
side first in the
+ // final order, matching a stable sort after rotating the compressed
left prefix.
+ if centroid_cmp(&left_centroid, &right_centroid) != Ordering::Less {
+ left_index -= 1;
+ left[output_index] = left_centroid;
+ } else {
+ right_index -= 1;
+ left[output_index] = right_centroid;
+ }
+ }
+ if right_index > 0 {
+ left[..right_index].copy_from_slice(&right[..right_index]);
+ }
+}
+
fn centroid_lower_bound(c: &Centroid, value: f64) -> Ordering {
if c.mean < value {
Ordering::Less
diff --git a/tests-integration/tests/serde_tests/tdigest.rs
b/tests-integration/tests/serde_tests/tdigest.rs
index 9702827..d955144 100644
--- a/tests-integration/tests/serde_tests/tdigest.rs
+++ b/tests-integration/tests/serde_tests/tdigest.rs
@@ -249,17 +249,28 @@ fn
test_serialized_bytes_stable_for_full_and_merged_digests() {
assert_eq!(bytes.len(), expected_len);
assert_eq!(fnv1a(&bytes), expected_hash);
}
+
+ let mut left = patterned_digest(10, 199, 2);
+ let left = left.serialize();
+ let mut left = TDigestMut::deserialize(&left, false).unwrap();
+ let mut right = patterned_digest(10, 199, 3);
+ let right = right.serialize();
+ let right = TDigestMut::deserialize(&right, false).unwrap();
+ left.merge(&right);
+ let bytes = left.serialize();
+ assert_eq!(bytes.len(), 272);
+ assert_eq!(fnv1a(&bytes), 0x5759_0428_c175_88ab);
}
#[test]
-fn test_updates_normalize_overfull_deserialized_staging_buffer() {
+fn test_updates_normalize_overfull_deserialized_buffer_without_centroids() {
let path = serialization_test_data("cpp_generated_files",
"tdigest_double_buf_n10_cpp.sk");
let mut bytes = fs::read(path).unwrap();
assert_eq!(&bytes[8..12], &0_u32.to_le_bytes()); // num centroids
assert_eq!(&bytes[12..16], &10_u32.to_le_bytes()); // num buffered
- // k=100 normally compresses at 840 buffered values. Extend a real C++
staging image just past
- // that producer threshold while keeping every added value within the
recorded min/max range.
+ // k=100 normally compresses at 840 buffered values. Extend a real C++
image without centroids
+ // just past that producer threshold while keeping every value within the
recorded min/max.
bytes[12..16].copy_from_slice(&841_u32.to_le_bytes());
for _ in 0..831 {
bytes.extend_from_slice(&10_f64.to_le_bytes());
@@ -273,7 +284,7 @@ fn
test_updates_normalize_overfull_deserialized_staging_buffer() {
assert_eq!(tdigest.total_weight(), 10_841);
assert_eq!(tdigest.min_value(), Some(1.0));
assert_eq!(tdigest.max_value(), Some(10.0));
- // The overfull image must not disable future compression and let the
staging buffer grow with
+ // The overfull image must not disable future compression and let the
buffered tail grow with
// every subsequent value.
assert!(tdigest.estimated_size() < 32_768);
let serialized = tdigest.serialize();
@@ -286,7 +297,7 @@ fn
test_updates_normalize_overfull_deserialized_staging_buffer() {
}
#[test]
-fn test_updates_normalize_overfull_deserialized_centroid_tail() {
+fn test_updates_normalize_overfull_deserialized_mixed_buffer() {
let path = serialization_test_data("cpp_generated_files",
"tdigest_double_buf_n1000_cpp.sk");
let mut bytes = fs::read(path).unwrap();
assert_eq!(&bytes[8..12], &89_u32.to_le_bytes()); // num centroids
diff --git a/tests-integration/tests/tdigest_test/sketch.rs
b/tests-integration/tests/tdigest_test/sketch.rs
index 162b6d5..5edc97f 100644
--- a/tests-integration/tests/tdigest_test/sketch.rs
+++ b/tests-integration/tests/tdigest_test/sketch.rs
@@ -79,7 +79,7 @@ fn test_maximum_k() {
}
#[test]
-fn test_estimated_size_releases_initial_staging_after_compression() {
+fn test_estimated_size_reuses_buffer_after_compression() {
const K: u16 = 200;
const TARGET_CENTROIDS: usize = 410;
const MAX_UNMERGED: usize = TARGET_CENTROIDS * 4;
@@ -94,9 +94,7 @@ fn
test_estimated_size_releases_initial_staging_after_compression() {
let size_before_compression = tdigest.estimated_size();
assert!(size_before_compression > inline_size);
tdigest.rank(0.5);
- // Unit-weight centroids are twice the size of staged f64 values. Allow
that representation
- // change while guarding against retaining both backing allocations after
compression.
- assert!(tdigest.estimated_size() <= size_before_compression * 2 +
inline_size);
+ assert!(tdigest.estimated_size() <= size_before_compression);
for value in MAX_UNMERGED..10_000 {
tdigest.update(value as f64);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]