kvstore: use batch delete to cleanup old events (#112737)
* use batchdelete for cleaning up old events * comment
This commit is contained in:
@@ -364,7 +364,7 @@ func (d *dataStore) Get(ctx context.Context, key DataKey) (io.ReadCloser, error)
|
||||
|
||||
// BatchGet retrieves multiple data objects in batches.
|
||||
// It returns an iterator that yields DataObj results for the given keys.
|
||||
// Keys are processed in batches (default 50) to balance between efficiency and memory usage.
|
||||
// Keys are processed in batches (default 50).
|
||||
// Non-existent entries will not appear in the result.
|
||||
func (d *dataStore) BatchGet(ctx context.Context, keys []DataKey) iter.Seq2[DataObj, error] {
|
||||
return func(yield func(DataObj, error) bool) {
|
||||
|
||||
@@ -15,7 +15,8 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
eventsSection = "unified/events"
|
||||
eventsSection = "unified/events"
|
||||
deleteEventBatchSize = 50
|
||||
)
|
||||
|
||||
// eventStore is a store for events.
|
||||
@@ -231,24 +232,43 @@ func (n *eventStore) ListSince(ctx context.Context, sinceRV int64) iter.Seq2[Eve
|
||||
|
||||
// CleanupOldEvents deletes events older than the specified retention period.
|
||||
func (n *eventStore) CleanupOldEvents(ctx context.Context, cutoff time.Time) (int, error) {
|
||||
deletedCount := 0
|
||||
|
||||
// Keys are stored in the format of "resource_version~namespace~group~resource~name"
|
||||
// With a start key of "1" and an end key of the cutoff time we can get all expired events.
|
||||
endKey := fmt.Sprintf("%d", snowflakeFromTime(cutoff))
|
||||
|
||||
// Collect keys to delete
|
||||
keysToDelete := make([]string, 0, deleteEventBatchSize)
|
||||
for key, err := range n.kv.Keys(ctx, eventsSection, ListOptions{StartKey: "1", EndKey: endKey}) {
|
||||
if err != nil {
|
||||
return deletedCount, fmt.Errorf("failed to list event keys: %w", err)
|
||||
return 0, fmt.Errorf("failed to list event keys: %w", err)
|
||||
}
|
||||
|
||||
// TODO should use batch deletes here when available
|
||||
if err := n.kv.Delete(ctx, eventsSection, key); err != nil {
|
||||
return deletedCount, fmt.Errorf("failed to delete event key %s: %w", key, err)
|
||||
}
|
||||
deletedCount++
|
||||
keysToDelete = append(keysToDelete, key)
|
||||
}
|
||||
|
||||
return deletedCount, nil
|
||||
// Use batch delete
|
||||
if err := n.batchDelete(ctx, keysToDelete); err != nil {
|
||||
return 0, fmt.Errorf("failed to batch delete events: %w", err)
|
||||
}
|
||||
|
||||
return len(keysToDelete), nil
|
||||
}
|
||||
|
||||
// batchDelete deletes multiple events in batches.
|
||||
// Keys are processed in batches (default 50).
|
||||
func (n *eventStore) batchDelete(ctx context.Context, keys []string) error {
|
||||
for len(keys) > 0 {
|
||||
batch := keys
|
||||
if len(batch) > deleteEventBatchSize {
|
||||
batch = batch[:deleteEventBatchSize]
|
||||
}
|
||||
keys = keys[len(batch):]
|
||||
|
||||
if err := n.kv.BatchDelete(ctx, eventsSection, batch); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// snowflake id with last two sections set to 0 (machine id and sequence)
|
||||
|
||||
@@ -610,3 +610,53 @@ func TestEventStore_CleanupOldEvents_EmptyStore(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 0, deletedCount, "Should not have deleted any events from empty store")
|
||||
}
|
||||
|
||||
func TestEventStore_BatchDelete(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
store := setupTestEventStore(t)
|
||||
|
||||
// Create multiple events (more than batch size to test batching)
|
||||
eventKeys := make([]string, 75)
|
||||
for i := 0; i < 75; i++ {
|
||||
event := Event{
|
||||
Namespace: "default",
|
||||
Group: "apps",
|
||||
Resource: "deployments",
|
||||
Name: "test-deployment",
|
||||
ResourceVersion: int64(1000 + i),
|
||||
Action: DataActionCreated,
|
||||
Folder: "test-folder",
|
||||
PreviousRV: int64(999 + i),
|
||||
}
|
||||
err := store.Save(ctx, event)
|
||||
require.NoError(t, err)
|
||||
|
||||
eventKeys[i] = EventKey{
|
||||
Namespace: event.Namespace,
|
||||
Group: event.Group,
|
||||
Resource: event.Resource,
|
||||
Name: event.Name,
|
||||
ResourceVersion: event.ResourceVersion,
|
||||
Action: event.Action,
|
||||
Folder: event.Folder,
|
||||
}.String()
|
||||
}
|
||||
|
||||
// Batch delete all events
|
||||
err := store.batchDelete(ctx, eventKeys)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Verify all events were deleted
|
||||
for i := 0; i < 75; i++ {
|
||||
_, err := store.Get(ctx, EventKey{
|
||||
Namespace: "default",
|
||||
Group: "apps",
|
||||
Resource: "deployments",
|
||||
Name: "test-deployment",
|
||||
ResourceVersion: int64(1000 + i),
|
||||
Action: DataActionCreated,
|
||||
Folder: "test-folder",
|
||||
})
|
||||
require.Error(t, err, "Event should have been deleted")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user