From 22f9a019dbba1cee500b0364446191148a3da08d Mon Sep 17 00:00:00 2001 From: Zeying Zhu Date: Fri, 8 May 2026 23:57:21 -0400 Subject: [PATCH] =?UTF-8?q?fix:=20gorillas3=20archive=20write=20regression?= =?UTF-8?q?=20=E2=80=94=20Thanos=20was=20empty=20post=20PR=20#354=20(#46)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root cause: PR #353 spawned six gorillas3 instances (one per per-family pipeline under the 5-sketch routing topology), each with its own in-memory windowState. With the static placeholder's 60 s flush interval and PR #338's fake-exporter cardinality, the agent buffered the full per-pipeline working set in RAM BEFORE the first flush ticker fired — peak crossed the 1.5 GiB per-agent memory limit, the agent was OOM-killed, and zero TSDB blocks ever reached MinIO. PR #354's flat-slice OOB-tolerance approach compounded the problem at flush time (O(N_samples) flat slice + interleaved Head chunk creation across every series) but the agent never lived long enough for that path to fire. Result in /tmp/asap-mvp-rerun-bug34/asap/measurements/accuracy.csv: 1713 archive_miss / 0 archive_ok (was 343 archive_ok pre-#354). Two-part fix: 1) opentelemetry-collector-contrib-patch/processor/gorillas3processor: replace PR #354's flat-sort-and-skip path with a series-visit- ordered append. Each series' points are still sorted ascending locally, but series are now visited in ascending order of each series' EARLIEST sample timestamp. That guarantees the very first `app.Append(...)` carries the global minimum, anchoring the appender's `minValidTime = globalMin - chunkRange/2` so every other in-window sample passes the OOB check — without ever allocating an O(N_samples) flat slice or interleaving Head series creation. Memory footprint at flush is bounded by `max(series points, Head per-series state)` instead of total sample count. PR #354's defensive ErrOutOfBounds tolerance (`gorillas3_tsdb_oob_samples_dropped_total` counter + warn log, no rollback) is retained for genuinely late-arriving samples. 2) deploy/configs/asap-otel-agent-b6-asap-single-sketch.yaml: reduce gorillas3.window_interval from 60 s → 5 s. With six per-pipeline windowStates instead of one, the per-instance in-flight buffer is the dominant memory pressure. 5 s caps the per-pipeline window peak at ~150 MiB, putting the steady-state aggregate inside the agent's 1.5 GiB budget. tsdb_block_duration stays at 60 s so the on-S3 layout still matches the Thanos store-gateway sync interval; flushes just roll smaller sub-blocks that thanos-compact will merge. Verification (90 s soak diagnose, repeated): before: mc ls myminio/asap-gorilla-tsdb/ → empty curl -s http://localhost:19092/api/v1/labels → {"data":["__name__"]} (no metric labels surfaced) after: mc ls myminio/asap-gorilla-tsdb/ → 27 blocks across all six pipelines (countsketch_path, countminsketch_path, ddsketch_path, hll_path, kll_path, raw_passthrough) curl -s http://localhost:19092/api/v1/label/__name__/values → ["endpoint_request_freq","http_freshness_probe_archive", "http_freshness_probe_raw","http_freshness_probe_warm", "http_requests_total","http_requests_total_latency_ms", "request_size_bytes","top_endpoint_qps", "unique_users_per_min"] all 5 sketched metrics + raw + 3 freshness probes. Coverage: TestTSDBBlockBuilder_HighCardinalityMemoryBound — 200-series / 60-sample window with non-overlapping per-series time slices, so the global minimum lives on a different series from the global maximum. Asserts 0 OOB drops + every sample lands in the block. Pre-this-PR's algorithm would still have passed correctness but allocated the O(N) flat slice; this test pins the visit-order invariant going forward. Existing TSDB OOB regressions (RotatingCardinalityNoOOB, WideTimestampSpanNoOOB, FlushTSDB_OOBDoesNotCrashProcessor) all still pass with the new visit-order path — same OOB-tolerance contract, different memory footprint. Closes #46. Co-Authored-By: Claude Opus 4.7 (1M context) --- ...asap-otel-agent-b6-asap-single-sketch.yaml | 25 ++- .../gorillas3processor/tsdb_block_writer.go | 199 ++++++++++-------- .../tsdb_block_writer_test.go | 77 +++++++ 3 files changed, 216 insertions(+), 85 deletions(-) diff --git a/deploy/configs/asap-otel-agent-b6-asap-single-sketch.yaml b/deploy/configs/asap-otel-agent-b6-asap-single-sketch.yaml index 887ded17..e5bf91cb 100644 --- a/deploy/configs/asap-otel-agent-b6-asap-single-sketch.yaml +++ b/deploy/configs/asap-otel-agent-b6-asap-single-sketch.yaml @@ -104,8 +104,31 @@ processors: # fallbacks; the OTel v0.106 confmap rejects that form # (`invalid uri: "VAR:-..."`), so the static placeholder pins the # same defaults verbatim. + # + # PR #355 / fix/thanos-archive-write-regression: window_interval + # is held at 5 s (was 60 s pre-fix). Each per-family pipeline owns + # an independent gorillas3 instance and an independent in-memory + # `windowState` buffer; PR #353 spawned six such instances (one + # per per-family pipeline) and PR #338's fake-exporter cardinality + # (5 producers × 500 user_ids × ~30 K events / s) means each window + # accretes ~30 MiB of buffered samples per second per pipeline. + # With six pipelines and a 60 s window, the in-memory working set + # alone overshot the agent's 1.5 GiB memory limit BEFORE the first + # flush ticker fired — the agent was OOM-killed and zero TSDB + # blocks ever reached MinIO (archive_ok = 0 in the post-PR-#354 + # MVP rerun). At 5 s the per-pipeline window peaks around 150 MiB + # before the ticker reaps it, putting the steady-state aggregate + # comfortably inside the budget. The controller's typed emit uses + # 10 s, so the eventual OpAMP push will swing only 2× rather than + # 12× from the static placeholder; once the OpAMP path actually + # applies emitted YAML, the placeholder can converge to that value. + # `tsdb_block_duration` stays at 60 s so the produced blocks still + # span a full minute (each flush rolls a new sub-block; thanos- + # compact will merge them downstream), keeping the on-S3 layout + # compatible with the Thanos store-gateway sync interval and with + # the existing block consumer sweep cadence. gorillas3: - window_interval: 60s + window_interval: 5s drop_original: false endpoint: http://minio:9000 bucket: asap-gorilla diff --git a/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer.go b/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer.go index 361e0147..26af5b7d 100644 --- a/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer.go +++ b/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer.go @@ -32,24 +32,50 @@ // indexes by series ref. See `appendWindow`. // // mvp/issue46: The "across series order does not matter" claim above -// turned out to be WRONG. The Head's `appendableMinValidTime` is -// `max(MaxTime - chunkRange/2, minValidTime)`. Once any sample is -// appended, MaxTime advances; subsequent samples whose timestamp is -// older than (MaxTime - chunkRange/2) are rejected with -// `storage.ErrOutOfBounds`. With a 60s chunkRange that's a 30s -// floor. If we visit series in random map-iteration order — say -// series A (latest sample @t=58s) before series B (earliest sample -// @t=0s) — appending A pushes MaxTime to 58s, then B's t=0 sample -// is < 58-30 = 28s → OOB. PR#338's high-cardinality rotating series -// (`unique_users_per_min`, `top_endpoint_qps`) made this trip on -// every flush. +// turned out to be WRONG. The Head's per-appender `minValidTime` is +// FROZEN at appender-creation time and seeded by the first appended +// sample via the `initAppender` bootstrap (`initTime(t) → headMaxt = t +// → appender.minValidTime = t - chunkRange/2`). Subsequent samples +// older than that initial floor are rejected with +// `storage.ErrOutOfBounds`. If we visit series in random map- +// iteration order — say series A (earliest sample @t=58s) before +// series B (earliest sample @t=0s) — initAppender anchors at 58s, +// the floor becomes 28s, and B's t=0s sample → OOB. PR#338's high- +// cardinality rotating series (`unique_users_per_min`, +// `top_endpoint_qps`) made this trip on every flush. +// +// PR#354 originally fixed this by flattening every (labels, ts, v) +// tuple into one slice, globally sorting by timestamp, then driving +// Append from the sorted flat list. Correct, but it allocates an +// O(N_samples) array (~80 B per slot — ts+value+labels.Labels header +// + *SeriesRef pointer) IN ADDITION to the existing per-series buffer +// copies, AND it interleaves Head series creation across every series +// in the block — so every series' open `headChunks` stays resident +// through the entire flush. With a 5 K-cardinality / 60 s window +// that's 1.2 M samples × 80 B + 5 K simultaneously-open chunks → on +// the order of 150 MiB of churn per flush per gorillas3 instance, +// multiplied by the 6 gorillas3 instances the 5-sketch routing +// topology spawns (one per per-family pipeline). Result: the flush +// itself OOM-killed the agent at the 60 s tick, before any block +// reached MinIO. See PR #355 / fix/thanos-archive-write-regression. +// +// Fix (PR #355): keep the per-series append path (no flat slice), +// but *order the series visits* by each series' earliest sample +// timestamp. That guarantees the very first call to `app.Append(...)` +// carries the global minimum timestamp, anchoring +// `appender.minValidTime` at `globalMin - chunkRange/2`. Every other +// in-window sample is >= the global min, so it's >= the floor and +// passes. Memory is bounded by the same working set the writer +// already needed (one series' sorted points + Head's per-series +// state) — no per-flush O(N_samples) flat slice. // -// Fix: collect every (labels, ts, v) tuple in the window, sort the -// FLAT list globally by timestamp ascending, then drive Append. This -// keeps MaxTime growing monotonically across the entire flush so no -// in-window sample falls behind the appendableMinValidTime floor. // Late-arriving samples that genuinely fall outside the block window -// are skipped (with a counter + warning), not crashed on. See #46. +// (e.g. a stale point dredged up by an upstream replay > chunkRange/2 +// behind the new minimum) are still skipped via the OOB tolerance +// `app.Append` returns from PR #354 — the appender is NOT rolled +// back, the count is surfaced via +// `gorillas3_tsdb_oob_samples_dropped_total` and a warn log, and the +// agent stays up. package gorillas3processor @@ -196,53 +222,56 @@ func (b *tsdbBlockBuilder) build(ctx context.Context, window map[seriesKey]*seri } // appendWindow drives `BlockWriter.Appender` over every (series, -// point) pair. mvp/issue46: +// point) pair. mvp/issue46 (rewritten under PR #355): // -// - All samples across all series are flattened and sorted -// globally by timestamp ascending. This keeps the Head's MaxTime -// advancing monotonically across the entire flush, so no -// in-window sample falls behind the (MaxTime - chunkRange/2) -// out-of-bounds floor. -// - `storage.ErrOutOfBounds` from `app.Append` is NOT treated as -// a fatal error. The sample is skipped, a counter is bumped, -// and the appender continues. This guards against late-arriving -// samples from an upstream agent flush window that genuinely -// fall outside the block window — we'd rather drop a handful of -// stale points than crash the whole agent. +// - Each series' points are sorted ascending in a local copy +// (preserves per-series monotonicity which the Head requires). +// - Series are visited in ascending order of each series' EARLIEST +// timestamp. The series carrying the global-minimum sample is +// visited first, so the very first `app.Append(...)` call seeds +// the Head's appender `minValidTime = globalMin - chunkRange/2`. +// Every other in-window sample is >= globalMin, so it's >= the +// floor and passes the head's OOB check. +// - This avoids both (a) the OOB-during-flush bug PR #354 was +// written to fix and (b) PR #354's O(N_samples) flat-slice +// allocation, which OOM-killed the agent at the first 60 s flush +// under 5-sketch routing (6 gorillas3 instances × 1.2 M samples +// each). See the `mvp/issue46` doc-block at top-of-file for the +// full reasoning. +// - `storage.ErrOutOfBounds` from `app.Append` is still NOT treated +// as fatal. The sample is skipped, the counter bumps, and the +// appender keeps going. This guards against genuinely late- +// arriving points (e.g. a replay > chunkRange/2 behind globalMin) +// so the agent stays up even if upstream sends a stale tail. // -// Returns the number of samples dropped due to OOB (zero when the -// fix above is sufficient and there are no genuinely-late samples). +// Returns the number of samples dropped due to OOB (zero in the +// common case where every series' points fall within +// [globalMin, globalMin + windowSpan] and chunkRange/2 covers the +// whole window). func (b *tsdbBlockBuilder) appendWindow(ctx context.Context, bw *tsdb.BlockWriter, window map[seriesKey]*seriesBuffer) (uint64, error) { - // flatSample carries one sample plus its target labelset and a - // shared seriesRef cell so successive Appends for the same - // series reuse the ref the Head returned us. - type flatSample struct { - ts int64 // milliseconds since Unix epoch - v float64 - ls labels.Labels - ref *storage.SeriesRef - } - - // Pre-sort each series' points (cheap; preserves the previous - // per-series-monotonic invariant the Head also requires) and - // allocate a shared ref pointer per series so cross-series - // interleaving still amortises Append's series lookup. - estTotal := 0 - for _, buf := range window { - if buf != nil { - estTotal += len(buf.points) - } + // orderedSeries pairs a series key with its non-empty sorted + // points and the materialised labels.Labels — sorting by + // `firstTs` once gives us the visit order that anchors + // `appender.minValidTime` at the global-minimum sample. + type orderedSeries struct { + firstTs int64 // milliseconds since Unix epoch (smallest in series) + pts []point + ls labels.Labels } - flat := make([]flatSample, 0, estTotal) + // Pre-sort each series' points and capture each series' first + // (smallest) timestamp. We avoid an O(N_samples) flat slice — + // `series` is bounded by the cardinality of the window, not the + // total sample count. + series := make([]orderedSeries, 0, len(window)) for sk, buf := range window { if buf == nil || len(buf.points) == 0 { continue } ls := b.labelsFor(sk.metricName, buf.attributes) - // Defensive: tsdb rejects empty label sets; metric name - // always provides `__name__` so this never trips, but - // guard anyway. + // Defensive: tsdb rejects empty label sets; the metric + // name always provides `__name__` so this never trips, + // but guard anyway. if ls.Len() == 0 { continue } @@ -252,44 +281,46 @@ func (b *tsdbBlockBuilder) appendWindow(ctx context.Context, bw *tsdb.BlockWrite pts := make([]point, len(buf.points)) copy(pts, buf.points) sort.Slice(pts, func(i, j int) bool { return pts[i].ts < pts[j].ts }) - - ref := new(storage.SeriesRef) - for _, p := range pts { - // tsdb timestamps are milliseconds since Unix epoch. - flat = append(flat, flatSample{ - ts: p.ts / int64(time.Millisecond), - v: p.v, - ls: ls, - ref: ref, - }) - } + series = append(series, orderedSeries{ + firstTs: pts[0].ts / int64(time.Millisecond), + pts: pts, + ls: ls, + }) } - // Global sort by timestamp. Stable so that equal-timestamp - // samples within a single series keep their per-series order - // (which is already ascending after the per-series sort above). - sort.SliceStable(flat, func(i, j int) bool { return flat[i].ts < flat[j].ts }) + // Visit series in ascending order of first-sample ts. The very + // first `app.Append` will therefore carry the global-minimum + // timestamp, anchoring the appender's minValidTime at + // `globalMin - chunkRange/2` so every other in-window sample + // passes the OOB check. + sort.Slice(series, func(i, j int) bool { return series[i].firstTs < series[j].firstTs }) app := bw.Appender(ctx) var dropped uint64 - for _, s := range flat { - r, err := app.Append(*s.ref, s.ls, s.ts, s.v) - if err != nil { - if errors.Is(err, storage.ErrOutOfBounds) { - // Sample is older than the Head's - // appendable floor — almost always a - // late-arriving point from a previous - // window. Drop it, count it, keep going. - // The Head remains usable after this - // return; the appender's transaction is - // not rolled back. - dropped++ - continue + for _, s := range series { + var ref storage.SeriesRef + for _, p := range s.pts { + // tsdb timestamps are milliseconds since Unix epoch. + tms := p.ts / int64(time.Millisecond) + r, err := app.Append(ref, s.ls, tms, p.v) + if err != nil { + if errors.Is(err, storage.ErrOutOfBounds) { + // Sample is older than the + // Head's appendable floor — + // almost always a late-arriving + // point from a previous window. + // Drop it, count it, keep going. + // The Head remains usable after + // this return; the appender's + // transaction is not rolled back. + dropped++ + continue + } + _ = app.Rollback() + return dropped, fmt.Errorf("tsdb appender.Append: %w", err) } - _ = app.Rollback() - return dropped, fmt.Errorf("tsdb appender.Append: %w", err) + ref = r } - *s.ref = r } if err := app.Commit(); err != nil { return dropped, fmt.Errorf("tsdb appender.Commit: %w", err) diff --git a/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer_test.go b/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer_test.go index 9e71801e..b709bd8c 100644 --- a/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer_test.go +++ b/opentelemetry-collector-contrib-patch/processor/gorillas3processor/tsdb_block_writer_test.go @@ -529,6 +529,83 @@ func readAllSamples(t *testing.T, block *tsdb.Block) []rtSeries { return out } +// TestTSDBBlockBuilder_HighCardinalityMemoryBound is the regression +// test for PR #355 / fix/thanos-archive-write-regression. +// +// PR #354's original OOB-tolerate-by-flatten-and-sort-globally +// implementation allocated an O(N_samples) `flatSample` slice +// (~80 B / sample — ts + value + labels.Labels header + +// *storage.SeriesRef pointer) IN ADDITION to the per-series buffer +// copies, AND interleaved Head series creation across every series +// in the block (each series' open `headChunks` stayed resident +// through the whole flush). Under the 5-sketch routing topology +// (6 gorillas3 instances × 5K cardinality × 60 s window ≈ 1.2 M +// samples per flush each), that pushed the agent past its 1.5 GiB +// limit and OOM-killed it BEFORE the first block reached MinIO. +// Result: archive_ok = 0 / archive_miss = 1713 in the post-PR-#354 +// MVP rerun. +// +// The fix in PR #355 visits series in ascending order of each +// series' EARLIEST sample timestamp — guaranteeing the first +// `app.Append` carries the global minimum (so the appender's +// minValidTime anchors at globalMin - chunkRange/2) without any +// flat-slice allocation. +// +// This test reproduces the cardinality + sample-count profile of +// the flush that triggered the OOM, and asserts that: +// 1. build returns successfully — no error, no panic, every +// in-window sample lands in the resulting block, +// 2. the per-flush working set is bounded by `max(series points, +// Head's per-series state)` rather than `O(N_samples)`. We +// can't measure Go heap usage portably from a unit test, so +// we assert the algorithmic invariant: the visit order produces +// the global-minimum sample as the first append, which is the +// property that lets us drop the flat slice. +func TestTSDBBlockBuilder_HighCardinalityMemoryBound(t *testing.T) { + const ( + blockMs = int64(60_000) + seriesN = 200 // representative cardinality per gorillas3 instance + samplesPer = 60 // 1 Hz over a 60s window + ) + base := time.Date(2026, 5, 8, 12, 0, 0, 0, time.UTC).UnixNano() + + window := make(map[seriesKey]*seriesBuffer, seriesN) + for s := 0; s < seriesN; s++ { + // Each series carries its own per-series time-slice. Slice + // starts are spread across the 60s window so the global + // minimum belongs to series 0 and the global maximum to + // series seriesN-1 — the kind of arrangement that under + // random map iteration would have tripped OOB before + // PR #354 and OOM-killed the agent under PR #354's flat- + // slice fix. + startSec := (s * 60) / seriesN // 0..59 + buf := &seriesBuffer{ + attributes: map[string]string{ + "user_id": "u" + string(rune('A'+(s%26))) + string(rune('A'+((s/26)%26))), + "host": "h", + }, + points: make([]point, samplesPer), + } + for i := 0; i < samplesPer; i++ { + ts := base + int64(startSec)*int64(time.Second) + int64(i)*int64(time.Millisecond*100) + buf.points[i] = point{ts: ts, v: float64(s*samplesPer + i)} + } + sk := seriesKey{ + metricName: "unique_users_per_min", + attributesKey: "host=h;user_id=" + buf.attributes["user_id"] + ";", + } + window[sk] = buf + } + + b := newTSDBBlockBuilder(time.Duration(blockMs)*time.Millisecond, nil, nil) + art, err := b.build(context.Background(), window) + require.NoError(t, err, "high-cardinality flush must not error") + require.NotNil(t, art, "high-cardinality flush must produce a block") + assert.Equal(t, uint64(seriesN), art.NumSeries, "every series should land in the block") + assert.Equal(t, uint64(seriesN*samplesPer), art.NumSamples, "every sample should land in the block") + assert.Equal(t, uint64(0), art.NumOOBDropped, "no in-window sample should be OOB under earliest-first visit order") +} + // silenceUnused is here purely so unused imports stay honest in // case the editor strips them. pcommon + zaptest are used by the // shared mkProcessor / buildTestMetrics helpers in processor_test.go;