search queue index mutex (#109474)
Use custom mutex for index, and don't hold it during BulkIndex.
This commit is contained in:
@@ -12,7 +12,6 @@ import (
|
||||
// It is responsible for ingesting events for a single index
|
||||
// It will batch events and send them to the index in a single bulk request
|
||||
type indexQueueProcessor struct {
|
||||
index ResourceIndex
|
||||
nsr NamespacedResource
|
||||
queue chan *WrittenEvent
|
||||
batchSize int
|
||||
@@ -20,8 +19,11 @@ type indexQueueProcessor struct {
|
||||
|
||||
resChan chan *IndexEvent // Channel to send results to the caller
|
||||
|
||||
mu sync.Mutex
|
||||
running bool
|
||||
indexMu sync.Mutex
|
||||
index ResourceIndex
|
||||
|
||||
runningMu sync.Mutex
|
||||
running bool
|
||||
}
|
||||
|
||||
type IndexEvent struct {
|
||||
@@ -34,8 +36,8 @@ type IndexEvent struct {
|
||||
}
|
||||
|
||||
func (b *indexQueueProcessor) updateIndex(newIndex ResourceIndex) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
b.indexMu.Lock()
|
||||
defer b.indexMu.Unlock()
|
||||
b.index = newIndex
|
||||
}
|
||||
|
||||
@@ -57,8 +59,8 @@ func (b *indexQueueProcessor) Add(evt *WrittenEvent) {
|
||||
b.queue <- evt
|
||||
|
||||
// Start the processor if it's not already running
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
b.runningMu.Lock()
|
||||
defer b.runningMu.Unlock()
|
||||
if !b.running {
|
||||
b.running = true
|
||||
go b.runProcessor()
|
||||
@@ -68,9 +70,9 @@ func (b *indexQueueProcessor) Add(evt *WrittenEvent) {
|
||||
// runProcessor is the task processing the queue of written events
|
||||
func (b *indexQueueProcessor) runProcessor() {
|
||||
defer func() {
|
||||
b.mu.Lock()
|
||||
b.runningMu.Lock()
|
||||
b.running = false
|
||||
b.mu.Unlock()
|
||||
b.runningMu.Unlock()
|
||||
}()
|
||||
|
||||
for {
|
||||
@@ -122,8 +124,9 @@ func (b *indexQueueProcessor) process(batch []*WrittenEvent) {
|
||||
} else {
|
||||
item.Action = ActionIndex
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
doc, err := b.builder.BuildDocument(ctx, evt.Key, evt.ResourceVersion, evt.Value)
|
||||
cancel()
|
||||
|
||||
if err != nil {
|
||||
result.Err = err
|
||||
} else {
|
||||
@@ -134,10 +137,11 @@ func (b *indexQueueProcessor) process(batch []*WrittenEvent) {
|
||||
req.Items = append(req.Items, item)
|
||||
}
|
||||
|
||||
b.mu.Lock()
|
||||
err := b.index.BulkIndex(req)
|
||||
b.mu.Unlock()
|
||||
b.indexMu.Lock()
|
||||
idx := b.index
|
||||
b.indexMu.Unlock()
|
||||
|
||||
err := idx.BulkIndex(req)
|
||||
if err != nil {
|
||||
for _, r := range resp {
|
||||
r.Err = err
|
||||
|
||||
Reference in New Issue
Block a user