Skip to content
Merged
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
14 changes: 11 additions & 3 deletions datasketches/src/frequencies/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,11 +131,19 @@ impl<T: Eq + Hash> FrequentItemsSketch<T> {
Self::with_lg_map_sizes(lg_max_map_size, LG_MIN_MAP_SIZE)
}

/// Returns true if the sketch is empty.
/// Returns true if the sketch has no active items.
///
/// A purge can remove all active items while retaining a non-zero total weight and
/// maximum error. Use [`Self::total_weight`] to distinguish that state from a newly created
/// or reset sketch.
pub fn is_empty(&self) -> bool {
self.hash_map.num_active() == 0
}

fn is_initial_state(&self) -> bool {
self.stream_weight == 0 && self.offset == 0 && self.hash_map.num_active() == 0
}

/// Returns the number of active items being tracked.
pub fn num_active_items(&self) -> usize {
self.hash_map.num_active()
Expand Down Expand Up @@ -359,7 +367,7 @@ impl<T: Eq + Hash> FrequentItemsSketch<T> {
where
T: Clone,
{
if other.is_empty() {
if other.is_initial_state() {
return;
}
let merged_total = self.stream_weight + other.stream_weight;
Expand Down Expand Up @@ -488,7 +496,7 @@ impl<T: Eq + Hash> FrequentItemsSketch<T> {
count_serialize_size: CountSerializeSize<T>,
serialize_item: SerializeItem<T>,
) -> Vec<u8> {
if self.is_empty() {
if self.is_initial_state() {
let mut bytes = SketchBytes::with_capacity(PREAMBLE_LONGS_EMPTY as usize * 8);
bytes.write_u8(PREAMBLE_LONGS_EMPTY);
bytes.write_u8(SERIAL_VERSION);
Expand Down
20 changes: 20 additions & 0 deletions datasketches/tests/frequencies_test/update.rs
Original file line number Diff line number Diff line change
Expand Up @@ -493,6 +493,26 @@ fn test_items_merge_empty_is_noop() {
assert_eq!(sketch.estimate(&1), 1);
}

#[test]
fn test_merge_preserves_purged_empty_state() {
let mut purged: FrequentItemsSketch<i64> = FrequentItemsSketch::new(32);
for item in 0..=(32 * 3 / 4) {
purged.update(item);
}
assert!(purged.is_empty());
assert_eq!(purged.total_weight(), 25);
assert_eq!(purged.maximum_error(), 1);

let mut merged: FrequentItemsSketch<i64> = FrequentItemsSketch::new(32);
merged.merge(&purged);

assert!(merged.is_empty());
assert_eq!(merged.num_active_items(), 0);
assert_eq!(merged.total_weight(), purged.total_weight());
assert_eq!(merged.maximum_error(), purged.maximum_error());
assert_eq!(merged.upper_bound(&1000), purged.upper_bound(&1000));
}

#[test]
fn test_row_equality_changes_with_updates() {
let mut sketch: FrequentItemsSketch<i32> = FrequentItemsSketch::new(8);
Expand Down
56 changes: 54 additions & 2 deletions datasketches/tests/serde_tests/frequencies.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,14 +104,66 @@ fn test_empty_round_trip() {
#[test]
fn test_purged_to_empty_round_trip() {
// Saturating the map with count-1 items makes the purge median 1, which
// removes every counter and leaves a non-trivial sketch empty.
// removes every counter while retaining stream and error state.
let mut sketch = FrequentItemsSketch::<i64>::new(32);
for i in 0..=(32 * 3 / 4) {
sketch.update(i);
}
assert!(sketch.is_empty());
let restored = FrequentItemsSketch::<i64>::deserialize(&sketch.serialize()).unwrap();
assert_eq!(sketch.num_active_items(), 0);
assert_eq!(sketch.total_weight(), 25);
assert_eq!(sketch.maximum_error(), 1);
assert_eq!(sketch.upper_bound(&1000), 1);

let bytes = sketch.serialize();
assert_eq!(bytes.len(), 4 * size_of::<u64>());
let restored = FrequentItemsSketch::<i64>::deserialize(&bytes).unwrap();
Comment thread
tisonkun marked this conversation as resolved.
assert!(restored.is_empty());
assert_eq!(restored.num_active_items(), 0);
assert_eq!(restored.total_weight(), sketch.total_weight());
assert_eq!(restored.maximum_error(), sketch.maximum_error());
assert_eq!(restored.upper_bound(&1000), sketch.upper_bound(&1000));
assert_eq!(restored.serialize(), bytes);
}

#[test]
fn test_zero_stream_weight_does_not_discard_other_state() {
// Simulate a wrapped stream weight or an inconsistent but accepted serialized image.
const STREAM_WEIGHT_OFFSET: usize = 2 * size_of::<u64>();

let mut active_sketch = FrequentItemsSketch::<i64>::new(32);
active_sketch.update_with_count(7, 3);
let mut active_bytes = active_sketch.serialize();
active_bytes[STREAM_WEIGHT_OFFSET..STREAM_WEIGHT_OFFSET + size_of::<u64>()].fill(0);

let active_restored = FrequentItemsSketch::<i64>::deserialize(&active_bytes).unwrap();
assert_eq!(active_restored.total_weight(), 0);
assert_eq!(active_restored.num_active_items(), 1);
assert_eq!(active_restored.estimate(&7), 3);
assert_eq!(active_restored.serialize(), active_bytes);

let mut active_merged = FrequentItemsSketch::<i64>::new(32);
active_merged.merge(&active_restored);
assert_eq!(active_merged.num_active_items(), 1);
assert_eq!(active_merged.estimate(&7), 3);

let mut purged_sketch = FrequentItemsSketch::<i64>::new(32);
for item in 0..=(32 * 3 / 4) {
purged_sketch.update(item);
}
let mut purged_bytes = purged_sketch.serialize();
purged_bytes[STREAM_WEIGHT_OFFSET..STREAM_WEIGHT_OFFSET + size_of::<u64>()].fill(0);

let purged_restored = FrequentItemsSketch::<i64>::deserialize(&purged_bytes).unwrap();
assert_eq!(purged_restored.total_weight(), 0);
assert_eq!(purged_restored.num_active_items(), 0);
assert_eq!(purged_restored.maximum_error(), 1);
assert_eq!(purged_restored.serialize(), purged_bytes);

let mut purged_merged = FrequentItemsSketch::<i64>::new(32);
purged_merged.merge(&purged_restored);
assert_eq!(purged_merged.num_active_items(), 0);
assert_eq!(purged_merged.maximum_error(), 1);
}

#[test]
Expand Down