implement bulkprocess in kv storage_backend

This commit is contained in:
Will Assis
2025-11-12 16:21:59 -03:00
parent cd15d73e32
commit e6b2a6cea3
4 changed files with 189 additions and 3 deletions
+1 -1
View File
@@ -338,7 +338,7 @@ type bulkRV struct {
counter int64
}
// When executing a bulk import we can fake the RV values
// Used when executing a bulk import where we can fake the RV values
func NewBulkRV() *bulkRV {
t := time.Now().Truncate(time.Second * 10)
return &bulkRV{
+1 -1
View File
@@ -449,7 +449,7 @@ func (d *dataStore) Delete(ctx context.Context, key DataKey) error {
return d.kv.Delete(ctx, dataSection, key.String())
}
func (n *dataStore) BatchDelete(ctx context.Context, keys []DataKey) error {
func (n *dataStore) batchDelete(ctx context.Context, keys []DataKey) error {
for len(keys) > 0 {
batch := keys
if len(batch) > dataBatchSize {
@@ -2971,7 +2971,7 @@ func TestDataStore_BatchDelete(t *testing.T) {
require.NoError(t, err)
}
err := ds.BatchDelete(ctx, keys)
err := ds.batchDelete(ctx, keys)
require.NoError(t, err)
// Verify all events were deleted
@@ -18,6 +18,7 @@ import (
"github.com/prometheus/client_golang/prometheus"
"go.opentelemetry.io/otel/trace"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"github.com/grafana/grafana/pkg/apimachinery/utils"
"github.com/grafana/grafana/pkg/storage/unified/resourcepb"
@@ -54,6 +55,7 @@ func convertEmptyToClusterNamespace(namespace string, withExperimentalClusterSco
type kvStorageBackend struct {
snowflake *snowflake.Node
kv KV
bulkLock *BulkLock
dataStore *dataStore
eventStore *eventStore
notifier *notifier
@@ -102,6 +104,7 @@ func NewKVStorageBackend(opts KVBackendOptions) (StorageBackend, error) {
backend := &kvStorageBackend{
kv: kv,
bulkLock: NewBulkLock(),
dataStore: newDataStore(kv),
eventStore: eventStore,
notifier: newNotifier(eventStore, notifierOptions{}),
@@ -1148,6 +1151,189 @@ func (k *kvStorageBackend) GetResourceLastImportTimes(ctx context.Context) iter.
}
}
func (b *kvStorageBackend) ProcessBulk(ctx context.Context, setting BulkSettings, iter BulkRequestIterator) *resourcepb.BulkResponse {
// TODO cross-node lock
err := b.bulkLock.Start(setting.Collection)
if err != nil {
return &resourcepb.BulkResponse{
Error: AsErrorResult(err),
}
}
defer b.bulkLock.Finish(setting.Collection)
bulkRvGenerator := NewBulkRV()
summaries := make(map[string]*resourcepb.BulkResponse_Summary, len(setting.Collection))
rsp := &resourcepb.BulkResponse{}
if setting.RebuildCollection {
for _, key := range setting.Collection {
events := make([]string, 0)
for evtKeyStr, err := range b.eventStore.ListKeysSince(ctx, 1) {
if err != nil {
b.log.Error("failed to list event: %s", err)
return rsp
}
evtKey, err := ParseEventKey(evtKeyStr)
if err != nil {
b.log.Error("error parsing event key: %s", err)
return rsp
}
if evtKey.Group != key.Group || evtKey.Resource != key.Resource || evtKey.Namespace != key.Namespace {
continue
}
events = append(events, evtKeyStr)
}
if err := b.eventStore.batchDelete(ctx, events); err != nil {
b.log.Error("failed to delete events: %s", err)
return rsp
}
historyKeys := make([]DataKey, 0)
for dataKey, err := range b.dataStore.Keys(ctx, ListRequestKey{
Namespace: key.Namespace,
Group: key.Group,
Resource: key.Resource,
}, SortOrderAsc) {
if err != nil {
b.log.Error("failed to list collection before delete: %s", err)
return rsp
}
historyKeys = append(historyKeys, dataKey)
}
previousCount := int64(len(historyKeys))
if err := b.dataStore.batchDelete(ctx, historyKeys); err != nil {
b.log.Error("failed to delete collection: %s", err)
return rsp
}
summaries[NSGR(key)] = &resourcepb.BulkResponse_Summary{
Namespace: key.Namespace,
Group: key.Group,
Resource: key.Resource,
PreviousCount: previousCount,
}
}
} else {
for _, key := range setting.Collection {
summaries[NSGR(key)] = &resourcepb.BulkResponse_Summary{
Namespace: key.Namespace,
Group: key.Group,
Resource: key.Resource,
}
}
}
obj := &unstructured.Unstructured{}
saved := make([]DataKey, 0)
rollback := func() {
// we don't have transactions in the kv store, so we simply delete everything we created
for _, val := range saved {
err = b.dataStore.Delete(ctx, val)
if err != nil {
b.log.Error("failed to delete during rollback: %s", err)
}
}
}
for iter.Next() {
if iter.RollbackRequested() {
rollback()
break
}
req := iter.Request()
if req == nil {
rollback()
rsp.Error = AsErrorResult(fmt.Errorf("missing request"))
break
}
rsp.Processed++
var action DataAction
switch resourcepb.WatchEvent_Type(req.Action) {
case resourcepb.WatchEvent_ADDED:
action = DataActionCreated
// Check if resource already exists for create operations
_, err := b.dataStore.GetLatestResourceKey(ctx, GetRequestKey{
Group: req.Key.Group,
Resource: req.Key.Resource,
Namespace: req.Key.Namespace,
Name: req.Key.Name,
})
if err == nil {
rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{
Key: req.Key,
Action: req.Action,
Error: "resource already exists",
})
continue
}
if !errors.Is(err, ErrNotFound) {
rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{
Key: req.Key,
Action: req.Action,
Error: fmt.Sprintf("failed to check if resource exists: %s", err),
})
continue
}
case resourcepb.WatchEvent_MODIFIED:
action = DataActionUpdated
case resourcepb.WatchEvent_DELETED:
action = DataActionDeleted
default:
rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{
Key: req.Key,
Action: req.Action,
Error: "invalid event type",
})
continue
}
err := obj.UnmarshalJSON(req.Value)
if err != nil {
rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{
Key: req.Key,
Action: req.Action,
Error: "unable to unmarshal json",
})
continue
}
dataKey := DataKey{
Group: req.Key.Group,
Resource: req.Key.Resource,
Namespace: req.Key.Namespace,
Name: req.Key.Name,
ResourceVersion: bulkRvGenerator.Next(obj),
Action: action,
Folder: req.Folder,
}
err = b.dataStore.Save(ctx, dataKey, bytes.NewReader(req.Value))
if err != nil {
rsp.Rejected = append(rsp.Rejected, &resourcepb.BulkResponse_Rejected{
Key: req.Key,
Action: req.Action,
Error: fmt.Sprintf("failed to save resource: %s", err),
})
continue
}
saved = append(saved, dataKey)
}
// TODO update last import time
return rsp
}
// readAndClose reads all data from a ReadCloser and ensures it's closed,
// combining any errors from both operations.
func readAndClose(r io.ReadCloser) ([]byte, error) {