diff --git a/frequencies/items_sketch.go b/frequencies/items_sketch.go index ee2c9c0..211b24b 100644 --- a/frequencies/items_sketch.go +++ b/frequencies/items_sketch.go @@ -471,14 +471,18 @@ func (i *ItemsSketch[C]) ToSlice() ([]byte, error) { preArr := make([]int64, preLongs) preArr[0] = pre0 preArr[1] = insertActiveItems(int64(activeItems), pre) - preArr[2] = int64(i.streamWeight) - preArr[3] = int64(i.offset) + preArr[2] = i.streamWeight + preArr[3] = i.offset for j := 0; j < preLongs; j++ { binary.LittleEndian.PutUint64(outArr[j<<3:], uint64(preArr[j])) } preBytes := preLongs << 3 - for j := 0; j < activeItems; j++ { - binary.LittleEndian.PutUint64(outArr[preBytes+j<<3:], uint64(i.hashMap.getActiveValues()[j])) + activeIndex := 0 + for slot, state := range i.hashMap.states { + if state > 0 { + binary.LittleEndian.PutUint64(outArr[preBytes+(activeIndex<<3):], uint64(i.hashMap.values[slot])) + activeIndex++ + } } copy(outArr[preBytes+(activeItems<<3):], bytes) } diff --git a/frequencies/items_sketch_test.go b/frequencies/items_sketch_test.go index d2f2ea8..d4d3db4 100644 --- a/frequencies/items_sketch_test.go +++ b/frequencies/items_sketch_test.go @@ -262,6 +262,66 @@ func TestSerializeDeserializeLong(t *testing.T) { assert.Equal(t, est, int64(1)) } +func TestItemsSketchToSlicePreservesValueOrder(t *testing.T) { + t.Run("int64", func(t *testing.T) { + testItemsSketchToSlicePreservesValueOrder( + t, + common.ItemSketchLongHasher{}, + common.ItemSketchLongSerDe{}, + func(index int) int64 { return int64(index) }, + ) + }) + t.Run("string", func(t *testing.T) { + testItemsSketchToSlicePreservesValueOrder( + t, + common.ItemSketchStringHasher{}, + common.ItemSketchStringSerDe{}, + func(index int) string { return "item-" + strconv.Itoa(index) }, + ) + }) +} + +func testItemsSketchToSlicePreservesValueOrder[C comparable]( + t *testing.T, + hasher common.ItemSketchHasher[C], + serde common.ItemSketchSerde[C], + itemAt func(int) C, +) { + const ( + mapSize = 64 + activeItems = mapSize * 3 / 4 + ) + + sketch, err := NewFrequencyItemsSketchWithMaxMapSize(mapSize, hasher, serde) + if !assert.NoError(t, err) { + return + } + for index := 0; index < activeItems; index++ { + err = sketch.UpdateMany(itemAt(index), int64(index+1)) + if !assert.NoError(t, err) { + return + } + } + + serialized, err := sketch.ToSlice() + if !assert.NoError(t, err) { + return + } + restored, err := NewFrequencyItemsSketchFromSlice(serialized, hasher, serde) + if !assert.NoError(t, err) { + return + } + + assert.Equal(t, sketch.GetNumActiveItems(), restored.GetNumActiveItems()) + assert.Equal(t, sketch.GetStreamLength(), restored.GetStreamLength()) + assert.Equal(t, sketch.GetMaximumError(), restored.GetMaximumError()) + for index := 0; index < activeItems; index++ { + actual, err := restored.GetLowerBound(itemAt(index)) + assert.NoError(t, err) + assert.Equal(t, int64(index+1), actual) + } +} + func TestResize(t *testing.T) { sketch1, err := NewFrequencyItemsSketchWithMaxMapSize[string](2<<_LG_MIN_MAP_SIZE, common.ItemSketchStringHasher{}, nil) for i := 0; i < 32; i++ { @@ -310,11 +370,10 @@ func TestMergeExact(t *testing.T) { assert.Equal(t, est, int64(1)) } -func TestNullMapReturns(t *testing.T) { +func TestEmptyMapReturnsNilActiveKeys(t *testing.T) { map1, err := newReversePurgeItemHashMap[int64](1<<_LG_MIN_MAP_SIZE, common.ItemSketchLongHasher{}, nil) assert.NoError(t, err) assert.Nil(t, map1.getActiveKeys()) - assert.Nil(t, map1.getActiveValues()) } func TestMisc(t *testing.T) { @@ -492,6 +551,65 @@ func BenchmarkItemSketch(b *testing.B) { } } +var benchmarkItemsSketchToSliceSink []byte + +func BenchmarkItemsSketchToSlice(b *testing.B) { + benchmarkItemsSketchToSlice( + b, + "int64", + common.ItemSketchLongHasher{}, + common.ItemSketchLongSerDe{}, + func(index int) int64 { return int64(index) }, + ) + benchmarkItemsSketchToSlice( + b, + "string", + common.ItemSketchStringHasher{}, + common.ItemSketchStringSerDe{}, + func(index int) string { return "item-" + strconv.Itoa(index) }, + ) +} + +func benchmarkItemsSketchToSlice[C comparable]( + b *testing.B, + itemType string, + hasher common.ItemSketchHasher[C], + serde common.ItemSketchSerde[C], + itemAt func(int) C, +) { + for _, mapSize := range []int{64, 256, 1024} { + activeItems := mapSize * 3 / 4 + b.Run(itemType+"/items="+strconv.Itoa(activeItems), func(b *testing.B) { + sketch, err := NewFrequencyItemsSketchWithMaxMapSize(mapSize, hasher, serde) + if err != nil { + b.Fatal(err) + } + for index := 0; index < activeItems; index++ { + if err := sketch.UpdateMany(itemAt(index), int64(index+1)); err != nil { + b.Fatal(err) + } + } + + serialized, err := sketch.ToSlice() + if err != nil { + b.Fatal(err) + } + benchmarkItemsSketchToSliceSink = serialized + b.SetBytes(int64(len(serialized))) + b.ReportAllocs() + b.ResetTimer() + + for iteration := 0; iteration < b.N; iteration++ { + serialized, err = sketch.ToSlice() + if err != nil { + b.Fatal(err) + } + benchmarkItemsSketchToSliceSink = serialized + } + }) + } +} + func generateTestRowItems(n int) []*RowItem[string] { items := make([]*RowItem[string], n) for i := 0; i < n; i++ { diff --git a/frequencies/reverse_purge_item_hash_map.go b/frequencies/reverse_purge_item_hash_map.go index 2e4f43b..fcc181a 100644 --- a/frequencies/reverse_purge_item_hash_map.go +++ b/frequencies/reverse_purge_item_hash_map.go @@ -249,19 +249,6 @@ func (r *reversePurgeItemHashMap[C]) hashDelete(deleteProbe int) { } } -func (r *reversePurgeItemHashMap[C]) getActiveValues() []int64 { - if r.numActive == 0 { - return nil - } - returnValues := make([]int64, 0, r.numActive) - for i := 0; i < len(r.values); i++ { - if r.states[i] > 0 { //isActive - returnValues = append(returnValues, r.values[i]) - } - } - return returnValues -} - func (r *reversePurgeItemHashMap[C]) getActiveKeys() []C { if r.numActive == 0 { return nil