From 2649c3e8391ca69aba9d1cf17954000b02f74073 Mon Sep 17 00:00:00 2001 From: Patrick Huie Date: Tue, 6 Oct 2026 13:16:26 -0400 Subject: [PATCH] adding batch size metric to durable emitter --- pkg/durableemitter/durable_emitter.go | 1 + pkg/durableemitter/durable_emitter_metrics.go | 20 +++++++ pkg/durableemitter/durable_emitter_test.go | 56 +++++++++++++++++++ 3 files changed, 77 insertions(+) diff --git a/pkg/durableemitter/durable_emitter.go b/pkg/durableemitter/durable_emitter.go index b075e71b23..225adb6c10 100644 --- a/pkg/durableemitter/durable_emitter.go +++ b/pkg/durableemitter/durable_emitter.go @@ -624,6 +624,7 @@ func (d *DurableEmitter) insertBatchLoop() { } ctx, cancel := context.WithTimeout(context.Background(), d.cfg.PublishTimeout) ids, batchErr := d.batchInserter.InsertBatch(ctx, payloads) + d.metrics.recordInsertBatchSize(ctx, len(payloads), batchErr) cancel() if batchErr == nil { d.eng.Debugw("DurableEmitter: coalesced insert flushed", "count", len(payloads)) diff --git a/pkg/durableemitter/durable_emitter_metrics.go b/pkg/durableemitter/durable_emitter_metrics.go index 2e1f4a2403..6883002428 100644 --- a/pkg/durableemitter/durable_emitter_metrics.go +++ b/pkg/durableemitter/durable_emitter_metrics.go @@ -56,6 +56,7 @@ type durableEmitterMetrics struct { expiredPurged metric.Int64Counter storeOps metric.Int64Counter storeOpDuration metric.Float64Histogram + insertBatchSize metric.Int64Histogram queueDepth metric.Int64Gauge queuePayloadBytes metric.Int64Gauge queueOldestAgeSec metric.Float64Gauge @@ -86,6 +87,10 @@ var durationBuckets = metric.WithExplicitBucketBoundaries( 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, ) +var batchSizeBuckets = metric.WithExplicitBucketBoundaries( + 1, 2, 4, 8, 16, 32, 64, 128, 256, 512, 1024, 2048, 4096, 8192, 16384, 32768, 65536, 131072, +) + // newDurableEmitterMetrics registers all DurableEmitter instruments on the // supplied meter. The caller is responsible for the meter's scope (the // instrument prefix below acts as the metric namespace). @@ -178,6 +183,14 @@ func newDurableEmitterMetrics(meter metric.Meter, clientName string) (*durableEm ); err != nil { return nil, err } + if m.insertBatchSize, err = meter.Int64Histogram( + "durable_emitter.insert_batch.size", + metric.WithUnit("{event}"), + metric.WithDescription("Events per coalesced InsertBatch flush to the durable store; labels: error={true,false}"), + batchSizeBuckets, + ); err != nil { + return nil, err + } if m.queueDepth, err = meter.Int64Gauge( "durable_emitter.queue.depth", metric.WithUnit("{row}"), @@ -289,6 +302,13 @@ func (m *durableEmitterMetrics) recordStoreOp(ctx context.Context, op string, el m.storeOpDuration.Record(ctx, elapsed.Seconds(), metric.WithAttributes(attribute.String("operation", op))) } +func (m *durableEmitterMetrics) recordInsertBatchSize(ctx context.Context, n int, batchErr error) { + if m == nil { + return + } + m.insertBatchSize.Record(ctx, int64(n), metric.WithAttributes(attribute.Bool("error", batchErr != nil))) +} + // recordQueueStats records the DB-derived queue statistics (payload bytes, // oldest pending age, TTL budget) from an already-observed snapshot. The // queue depth gauge itself is recorded separately by DurableEmitter from the diff --git a/pkg/durableemitter/durable_emitter_test.go b/pkg/durableemitter/durable_emitter_test.go index 75329a4427..9d1af319be 100644 --- a/pkg/durableemitter/durable_emitter_test.go +++ b/pkg/durableemitter/durable_emitter_test.go @@ -815,6 +815,53 @@ func TestDurableEmitter_MetricsRegistersQueueTTLBudget(t *testing.T) { }, 2*time.Second, 10*time.Millisecond, "expected durable_emitter.queue.ttl_budget_seconds Int64 gauge in exported metrics") } +func TestDurableEmitter_MetricsRecordsInsertBatchSize(t *testing.T) { + meter, reader := newTestMeter(t) + + store := NewMemDurableEventStore() + be := newTestBatchEmitter() + cfg := DefaultConfig() + cfg.RetransmitInterval = time.Hour + cfg.InsertBatchSize = 10 + cfg.InsertBatchWorkers = 1 + cfg.InsertBatchFlushInterval = 100 * time.Millisecond + cfg.Metrics = &DurableEmitterMetricsConfig{PollInterval: 25 * time.Millisecond} + + em, err := NewDurableEmitter(store, be, true, cfg, logger.Test(t), meter) + require.NoError(t, err) + servicetest.Run(t, em) + ctx := t.Context() + + const n = 25 + var wg sync.WaitGroup + for range n { + wg.Go(func() { assert.NoError(t, em.Emit(ctx, []byte("m"), testEmitAttrs()...)) }) + } + wg.Wait() + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(ctx, &rm)) + + var found bool + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name != "durable_emitter.insert_batch.size" { + continue + } + h, ok := m.Data.(metricdata.Histogram[int64]) + require.True(t, ok, "expected Histogram[int64] for %s, got %T", m.Name, m.Data) + for _, dp := range h.DataPoints { + if hasMetricBoolAttr(dp.Attributes, "error", false) { + found = true + assert.Equal(t, int64(n), dp.Sum) + assert.Less(t, dp.Count, uint64(n), "batches should be coalesced, not one flush per event") + } + } + } + } + assert.True(t, found, "expected durable_emitter.insert_batch.size in exported metrics") +} + func counterSumByPhase(t *testing.T, rm metricdata.ResourceMetrics, name, phase string) int64 { t.Helper() var total int64 @@ -882,6 +929,15 @@ func hasMetricStringAttr(set attribute.Set, key, want string) bool { return false } +func hasMetricBoolAttr(set attribute.Set, key string, want bool) bool { + for _, kv := range set.ToSlice() { + if string(kv.Key) == key { + return kv.Value.AsBool() == want + } + } + return false +} + func TestDurableEmitter_MetricsPublishBatchEventPhase(t *testing.T) { meter, reader := newTestMeter(t)