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
171 changes: 171 additions & 0 deletions internal/telemetry/commit_memory_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,171 @@
package telemetry

import (
"context"
"fmt"
"runtime"
"runtime/debug"
"sync"
"sync/atomic"
"testing"
"time"
)

// commitPeakRatioLimit is how much live heap one CommitBatch may hold relative
// to the payload it was handed. The number that matters is the multiplier, not
// an absolute byte count: peak scales with batch size, and batch size scales
// with ingest rate, so a process sized for a quiet hour is sized wrong for a
// busy one.
//
// This is a ratchet, not a target, and it is placed between two measured
// distributions rather than guessed. Under -race, dictionary-encoded
// near-unique JSON columns ran 4.09x-4.92x across repeated runs; writing them
// plain runs 2.79x-3.49x. The limit sits in the gap, so it passes the current
// encoding with headroom for GC timing noise and fails a regression to the one
// that was OOM-killing the process.
//
// What remains above 1x is a conversion copy of every JSON column
// (makeSpanParquetRow) plus ~24 MiB of row-struct headers per 50k spans and
// the writer's own buffers. Lower this constant as each of those is addressed;
// see the ordered plan in the ingest-memory issue.
const commitPeakRatioLimit = 3.75

// syntheticSpan builds a span whose JSON columns look like production: a
// resource block shared across the batch, and attributes/events/links that are
// near-unique per row. That distinction is the whole point — a dictionary over
// a near-unique column stores every value twice.
func syntheticSpan(i int, base time.Time) Span {
start := base.Add(time.Duration(i) * time.Microsecond)
return Span{
Namespace: "prod",
TraceID: fmt.Sprintf("%032x", i),
SpanID: fmt.Sprintf("%016x", i),
ServiceName: "checkout",
Name: "POST /api/orders",
Kind: "SERVER",
StartUnixNanos: start.UnixNano(),
EndUnixNanos: start.Add(3 * time.Millisecond).UnixNano(),
DurationMS: 3,
StatusCode: "OK",
ResourceJSON: []byte(`{"service.name":"checkout","deployment.environment":"prod","host.name":"cube-10"}`),
AttributesJSON: []byte(fmt.Sprintf(
`{"http.method":"POST","http.route":"/api/orders","http.status_code":200,"order.id":"%d","user.id":"u-%d","session":"%032x"}`, i, i*7, i)),
EventsJSON: []byte(fmt.Sprintf(
`[{"name":"validated","time":%d},{"name":"charged","time":%d,"amount":%d}]`, start.UnixNano(), start.UnixNano()+1000, i%9999)),
LinksJSON: []byte(fmt.Sprintf(`[{"trace_id":"%032x","span_id":"%016x"}]`, i+1, i+1)),
}
}

func payloadBytes(spans []Span) int64 {
var total int64
for i := range spans {
total += int64(len(spans[i].ResourceJSON) + len(spans[i].AttributesJSON) +
len(spans[i].EventsJSON) + len(spans[i].LinksJSON) +
len(spans[i].TraceID) + len(spans[i].SpanID) + len(spans[i].ServiceName) + len(spans[i].Name))
}
return total
}

// samplePeakHeapInuse runs fn while polling HeapInuse, and reports the highest
// value observed above the pre-run baseline. GOGC is pinned low for the
// duration so the figure approximates live bytes rather than live bytes plus
// whatever slack the collector happened to be carrying; production GOGC roughly
// doubles it.
func samplePeakHeapInuse(fn func()) int64 {
previousGC := debug.SetGCPercent(5)
defer debug.SetGCPercent(previousGC)

runtime.GC()
var stats runtime.MemStats
runtime.ReadMemStats(&stats)
baseline := int64(stats.HeapInuse)

var peak atomic.Int64
stop := make(chan struct{})
var sampler sync.WaitGroup
sampler.Add(1)
go func() {
defer sampler.Done()
var sample runtime.MemStats
for {
select {
case <-stop:
return
default:
}
runtime.ReadMemStats(&sample)
if delta := int64(sample.HeapInuse) - baseline; delta > peak.Load() {
peak.Store(delta)
}
time.Sleep(2 * time.Millisecond)
}
}()

fn()
close(stop)
sampler.Wait()
return peak.Load()
}

// One commit must not cost several times the payload it was given. Four commit
// workers run concurrently, each queued behind more batches, so a multiplier
// here is a multiplier on the whole process.
func TestCommitBatchPeakHeapStaysNearPayloadSize(t *testing.T) {
if testing.Short() {
t.Skip("allocates ~50k spans")
}
store, err := OpenParquetStore(t.TempDir())
if err != nil {
t.Fatalf("open parquet store: %v", err)
}
base := time.Date(2026, 9, 20, 12, 0, 0, 0, time.UTC)
spans := make([]Span, 50_000)
for i := range spans {
spans[i] = syntheticSpan(i, base)
}
raw := payloadBytes(spans)

peak := samplePeakHeapInuse(func() {
if err := store.CommitBatch(context.Background(), BatchMetadata{ID: "peak-heap"}, spans, nil, nil); err != nil {
t.Errorf("commit batch: %v", err)
}
})
runtime.KeepAlive(spans)

ratio := float64(peak) / float64(raw)
t.Logf("payload %.1f MiB, peak heap %.1f MiB, ratio %.2fx",
float64(raw)/(1<<20), float64(peak)/(1<<20), ratio)
if ratio > commitPeakRatioLimit {
t.Errorf("CommitBatch peak heap %.2fx payload, limit %.2fx: one batch holds several copies of what it was handed",
ratio, commitPeakRatioLimit)
}
}

// BenchmarkCommitBatch reports peak heap alongside time so a regression shows
// up as a number rather than as an out-of-memory kill in production.
func BenchmarkCommitBatch(b *testing.B) {
base := time.Date(2026, 9, 20, 12, 0, 0, 0, time.UTC)
spans := make([]Span, 50_000)
for i := range spans {
spans[i] = syntheticSpan(i, base)
}
raw := payloadBytes(spans)

var peak int64
for i := 0; b.Loop(); i++ {
store, err := OpenParquetStore(b.TempDir())
if err != nil {
b.Fatalf("open parquet store: %v", err)
}
observed := samplePeakHeapInuse(func() {
if err := store.CommitBatch(context.Background(), BatchMetadata{ID: fmt.Sprintf("bench-%d", i)}, spans, nil, nil); err != nil {
b.Fatalf("commit batch: %v", err)
}
})
if observed > peak {
peak = observed
}
}
b.ReportMetric(float64(peak)/(1<<20), "peak-heap-MiB")
b.ReportMetric(float64(peak)/float64(raw), "peak/payload")
}
25 changes: 18 additions & 7 deletions internal/telemetry/parquet_rows.go
Original file line number Diff line number Diff line change
@@ -1,5 +1,16 @@
package telemetry

// Dictionary encoding is applied per column chunk and holds every distinct
// value for the whole row group, which ingest flushes only at
// parquetRowGroupRows. That pays for itself on a column whose values repeat —
// service names, HTTP methods, scope identifiers, the resource block shared by
// every span in a ResourceSpans — and costs a full second copy of the payload
// on a column whose values do not. The JSON columns carrying per-row
// attributes, events, links and exemplars are near-unique by construction, so
// they are written plain: dictionary-encoding them made one 50k-span commit
// hold several times the batch it was handed, which is what
// TestCommitBatchPeakHeapStaysNearPayloadSize pins.

type spanParquetRow struct {
Namespace string `parquet:"namespace"`
TraceID string `parquet:"trace_id"`
Expand All @@ -17,9 +28,9 @@ type spanParquetRow struct {
Status string `parquet:"status"`
StatusMessage string `parquet:"status_message"`
ResourceJSON string `parquet:"resource_json,dict"`
AttributesJSON string `parquet:"attributes_json,dict"`
EventsJSON string `parquet:"events_json,dict"`
LinksJSON string `parquet:"links_json,dict"`
AttributesJSON string `parquet:"attributes_json"`
EventsJSON string `parquet:"events_json"`
LinksJSON string `parquet:"links_json"`
TraceState string `parquet:"trace_state,dict"`
Flags int64 `parquet:"flags"`
ScopeName string `parquet:"scope_name,dict"`
Expand All @@ -36,7 +47,7 @@ type spanParquetRow struct {
ServiceVersion string `parquet:"service_version,dict"`
DeploymentEnv string `parquet:"deployment_env,dict"`
ExceptionType string `parquet:"exception_type,dict"`
ExceptionMessage string `parquet:"exception_message,dict"`
ExceptionMessage string `parquet:"exception_message"`
}

func makeSpanParquetRow(r Span) spanParquetRow {
Expand Down Expand Up @@ -90,7 +101,7 @@ type logParquetRow struct {
SpanID string `parquet:"span_id"`
Flags int64 `parquet:"flags"`
ResourceJSON string `parquet:"resource_json,dict"`
AttributesJSON string `parquet:"attributes_json,dict"`
AttributesJSON string `parquet:"attributes_json"`
ScopeName string `parquet:"scope_name,dict"`
ScopeVersion string `parquet:"scope_version,dict"`
IngestedAt int64 `parquet:"ingested_at,timestamp(nanosecond)"`
Expand Down Expand Up @@ -123,8 +134,8 @@ type metricParquetRow struct {
HistCountsJSON string `parquet:"hist_counts_json,dict"`
HistCount int64 `parquet:"hist_count"`
HistSum float64 `parquet:"hist_sum"`
ExemplarsJSON string `parquet:"exemplars_json,dict"`
AttributesJSON string `parquet:"attributes_json,dict"`
ExemplarsJSON string `parquet:"exemplars_json"`
AttributesJSON string `parquet:"attributes_json"`
ResourceJSON string `parquet:"resource_json,dict"`
ScopeName string `parquet:"scope_name,dict"`
ScopeVersion string `parquet:"scope_version,dict"`
Expand Down
Loading