184
185
// bulkBeforeAction is triggered before bulk processor commit
186
>
func (p *processorImpl) bulkBeforeAction(_ int64, requests []elastic.BulkableRequest) {
processor.go
187
>
metrics.ElasticsearchBulkProcessorRequests.With(p.metricsHandler).Record(int64(len(requests)))
188
>
p.metricsHandler.Histogram(metrics.ElasticsearchBulkProcessorBulkSize.Name(), metrics.ElasticsearchBulkProcessorBulkSize.Unit()).
189
>
Record(int64(len(requests)))
190
>
191
>
for _, request := range requests {
192
>
visibilityTaskKey := p.extractVisibilityTaskKey(request)
193
>
if visibilityTaskKey == "" {
194
continue
195
}
196
>
_, _, _ = p.mapToAckFuture.GetAndDo(visibilityTaskKey, func(key any, value any) error {
processor.go
197
>
ackF, ok := value.(*ackFuture)
198
>
if !ok {
199
p.logger.Fatal(fmt.Sprintf("mapToAckFuture has item of a wrong type %T (%T expected).", value, &ackFuture{}), tag.Value(key))
200
}
202
>
return nil
203
})
204
}