Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions pkg/durableemitter/durable_emitter.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
20 changes: 20 additions & 0 deletions pkg/durableemitter/durable_emitter_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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}"),
Expand Down Expand Up @@ -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
Expand Down
56 changes: 56 additions & 0 deletions pkg/durableemitter/durable_emitter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)

Expand Down
Loading