Skip to content

Correlate large captures in a streaming columnar pass - #1

Merged
lhotari merged 22 commits into
mainfrom
correlator-streaming-join
Sep 22, 2026
Merged

lhotari merged 22 commits into
mainfrom
correlator-streaming-join

Conversation

@lhotari

@lhotari lhotari commented Sep 22, 2026

Copy link
Copy Markdown
Collaborator

Motivation

The offline correlator could not correlate a real capture.

A Pulsar broker run produced 1,123,638 selected intervals — a 110.9 MB capture stream and a 38.2 MB JFR carrying 1,121,418 profiler.SignalSample events over 10,631 distinct Java stacks. Correlating it exhausted an 8 GB heap outright. At 14 GB it collapsed the stacks in 51 seconds and then died with Synthetic event limit exceeded while writing the synthetic JFR — after all the work was done.

Two independent causes:

  1. Every decoded record became a retained Gson DOM. CaptureInput read the whole capture into List<JsonObject> and Map<Long, JsonArray> before the join started, so the footprint scaled with the number of recorded intervals at roughly 8 KB per interval rather than tens of bytes.
  2. Nothing was written until everything was accumulated. The synthetic-JFR writer discovered its event-count limit mid-write, so a job that was going to fail still paid for the entire correlation first.

Design

A columnar two-pass architecture.

Pass 1 streams both inputs into parallel primitive arrays — one growable column per field, ~38 bytes per row — joins them on the exact 64-bit cookie through an open-addressed long -> int index, and accumulates collapsed weights. No DOM is retained. JFR stacks and threads are interned by folding a 64-bit hash over the raw frame list and confirming a hit against the stored frames, so a frame array is allocated only for a genuinely new stack.

Pass 2 re-reads both digest-verified files to emit the per-row audit outputs through a Sink. The CLI's sink writes JSONL as rows stream past; the library's collects into lists and rebuilds today's Analysis exactly. That is what lets the existing fixtures pass unchanged while the CLI never materialises anything.

Quantum pre-planning. The number of synthetic events a quantum produces is exactly events(q) = sum over stacks of floor(W_s / q), derivable from the per-stack weights before a single event is written. The quantum is now raised until the count fits and recorded in the report, so the write no longer fails after the work is done.

A degradation ladder. Integrity failures stay fatal; volume problems degrade; data-quality problems are reported. The rungs are: drop audit outputs, thin the source, narrow the window. Thinning is deterministic — the keep/drop test hashes the cookie, which is the join key, so an observation and its sample are dropped together and the counts keep describing one coherent subsample — and is reweighted by the exact reciprocal of the realised probability, composing with the kernel's existing admissionThreshold inverse-probability estimation rather than replacing it.

Naming distinguishes the two degradations, because they mean different things: thinning yields a statistically valid estimate of the whole window (ordinary names, exit 0); window narrowing yields a result covering less than what was asked for (INCOMPLETE-jonoffcpu-*, a jonoffcpu-narrowed.json marker, no jonoffcpu-complete.json, exit 2).

Arithmetic. Rather than a dual BigInteger/long path, one rule: monotonic timestamps and the clock offset must fit a signed 64-bit value (2^63 ns is 292 years), which makes plain Math.*Exact safe. sourceAggregate keeps exact BigInteger fixed-point.

Results

Against the capture above:

Heap Result Wall Peak RSS Peak retention
-Xmx14g (reference) exit 0 2:00 2.96 GB 203 MiB
-Xmx2g exit 0 2:02 1.28 GiB 203 MiB
-Xmx512m exit 0 1:20 723 MB 203 MiB

890,086 matches at every heap size, with SHA256-identical collapsed stacks. Measured retention is ~190 bytes per recorded interval. A synthetic 2,000,000-observation fixture correlates within -Xmx1g on every check.

What does not change

The capture format, the agent, and docs/schema/jonoffcpu-capture.proto are untouched. Every validation rule and failure message still fires on the same input with the same text, bar the deviations listed below. The public API — OffCpuCorrelator.correlate(source, jfr, output), run(args), and the shapes of Analysis, PartialAnalysis, Match, ClassifiedRecord, PopulationEstimate, Limits — is unchanged, and OfflineCorrelatorTest, PartialCorrelatorTest and CompatibilityJfrWriterTest pass with no source changes. Report schemaVersion stays 1. Synthetic JFR output stays byte-identical apart from the quantum.

Behaviour changes worth review

  • --audit defaults to matches on the CLI (the library default stays FULL), so jonoffcpu-classified-records.jsonl is no longer written by default. This is backward-incompatible; both privileged Python proof tools that read it now pass --audit full, and the docs call it out.
  • Limits rejects a time boundary outside the signed 64-bit nanosecond range with Time boundary outside signed 64-bit nanoseconds. Values in [2^63, 2^64) that previously constructed are now rejected.
  • A handler delay in (Long.MAX_VALUE, U64_MAX] now classifies as invalid-handler-delay rather than handler-delay-limit-exceeded. Requires an absurd clock offset.
  • New options: --audit, --thinning, --thinning-seed, --on-limit, plus a degradation object in the report.

Known limitations

  • Spec section 2 (the windowed merge join with straggler spill) is not implemented, so retention still scales with capture length and its acceptance criterion is knowingly unmet. It was gated on this measurement and the measurement said it was not needed for captures of this size — the real capture fits in 512 MB. The work is self-contained and additive if longer captures later need it.
  • Below ~512 MB, three OOMs escape the ladder — in the synthetic-JFR writer (JMC's LEB128ByteArrayWriter), in the JDK's own JFR ChunkParser during the second read pass, and in the cookie index constructor. The ladder governs the retention budget it accounts for; these are third-party allocations it does not. Well below the 2 GB target, but a real gap.
  • KnownWaitAttributionCheck compiles against the rebuilt façade but was not executed — it needs a privileged BPF host. Its runtime behaviour wants CI verification.

Verification

./gradlew :jonoffcpu-correlator:check --no-daemon is green at every one of the 22 commits. New fixtures cover the primitive structures, golden equivalence between the streamed and retained paths, audit levels, synthetic-JFR ordering and timestamps, the 2M-row scale case, thinning accuracy and determinism, the degradation ladder, and the audit re-read under a degraded run.

The correlation join keys on a 64-bit cookie; a HashMap<String, List<Row>> spends
about 150 bytes an entry on what an open-addressed long -> int map holds in 16.
U64 keeps the u64 semantics the capture format guarantees while the columns store
raw bits, and rejects a monotonic timestamp above 2^63 by name rather than wrapping.
The capture stream interns its native stacks; the JFR does not, and the correlator
expanded a copy into every sample. A three-minute broker capture carries 10,631
distinct stacks across 1.12 million samples, so the expansion is about 100x
redundant and is the dominant term in the correlator's footprint. Interning is
done without building a key per sample: a hash is folded over the raw frame list
and confirmed against the stored frames, so only a new stack allocates.
canonicalStack collapses a missing key, a JSON null, and an empty string to the
same segment for every frame field, and a missing numeric field to the same
segment too; JfrDictionaries distinguishes null from empty and (in principle) a
missing number from an explicit zero. Document why that divergence is
unreachable in practice (SignalJfrExporter.frame is the sole producer and
always populates all six fields together) and why, where it could differ, the
interner is strictly finer than canonicalStack rather than coarser -- it can
only split what should merge, never merge what should stay distinct. No
behavior change.
Correlation's footprint scaled with the number of recorded intervals because every
record stayed a Gson object until the analysis was returned. These columns hold the
38 bytes a row the join and the aggregates actually read, and give the
classification vocabulary a single definition instead of scattered string literals.
CaptureInput now hands each observation and stack to a visitor as it reads them and
keeps only the announced stack ids, so the reader's footprint stops scaling with the
number of recorded intervals. The admission budget is split accordingly: a row
counter over both inputs, the retained control documents, and the live size of the
engine's primitive structures, which is what --max-retained-bytes was always
supposed to be guarding.

OfflineCorrelator passes a visitor that rebuilds today's observation and stack lists
locally, so its join and validation logic are unchanged; PartialCorrelatorTest's
direct CaptureInput.readPartial call updates to the new visitor parameter.
CaptureInput.SourceVisitor.stack now hands the decoded CaptureProto.Stack to the
visitor instead of just its frame count, so OfflineCorrelator's interim RowCollector
can rebuild today's stack-id -> frames map with real address/symbol/module content
via CaptureStream.stackRow, matching the previous retained behaviour exactly. Task A4
must leave observable behaviour unchanged even though it is committed standalone, and
placeholder frames in jonoffcpu-classified-records.jsonl would have been a silent
regression until Task A6 rebuilt the audit pass.
CorrelationEngine streams a finalized capture and its combined JFR into
SourceColumns/SampleColumns/JfrDictionaries, joins them on the exact
64-bit cookie through LongIntMap, and accumulates duration-weighted
stacks per interned stack id, reproducing the retained chain's reason
precedence, three-phase duplicate invalidation, post-join clipping and
row numbering exactly. CorrelationResult is its output: columns,
dictionaries, cookie index and counters, with per-match values (clipped
interval, delivery delay) derived rather than stored.

The engine has no caller yet: OfflineCorrelator.correlate/correlatePartial
still run the old retained-DOM join (validateSource/join/sourceAggregate/
populationEstimate/Row), which Task A6 replaces and then deletes.
OfflineCorrelator.readJfr and validateStats are extracted to
package-private statics (behavior unchanged) so the engine can drive the
same JFR dispatch.

Limits' compact constructor now rejects a time boundary at or above
2^63 ns instead of 2^64-1 ("Time boundary outside signed 64-bit
nanoseconds"): the engine's columns hold clock values as signed longs,
and 2^63 ns is 292 years of uptime, so this only rejects a corrupt
input. This is one of the plan's four deliberate deviations.

Adds a StreamingCorrelatorTest golden-equivalence fixture verifying the
CLI (--audit full) and the library facade (OffCpuCorrelator.write)
produce byte-identical collapsed/classified-records/matches output and
equal reports. Exercising it required a minimal --audit
none|matches|full CLI flag ahead of its planned task slot (A8); the CLI
default stays "full" and OutputOptions.defaults() is unaffected, so
this only gates whether the two per-row audit files are written.
OfflineCorrelator.correlate/correlatePartial now call CorrelationEngine
directly and rebuild today's in-memory Analysis/PartialAnalysis through a
new AuditPass: a second, re-reading pass over the already digest-verified
capture and JFR that recomputes each row's classification from the columnar
join and hands it to a Sink. The CLI's audit files come from the same pass;
the library's collecting Sink (OfflineCorrelator.Collector) rebuilds the
classified-records list and the source-ordered match array so
OfflineCorrelatorTest, PartialCorrelatorTest and KnownWaitAttributionCheck
keep consuming analysis.records()/matches() unchanged.

The old retained-DOM path (validateSource, join, invalidateDuplicates,
sourceAggregate, populationEstimate, optionalCounter, index, stack, Row,
RowCollector) is deleted from OfflineCorrelator; CorrelationEngine already
carried its own copies of the join and the population estimate. Added
CorrelationResult.sourceIndex so a matched JFR sample's cookie can resolve
back to its source slot without a retained per-row map.

Added AuditLevel as the brief specifies; nothing in this commit constructs
one yet, since the CLI still uses the string-based --audit stopgap from
Task A5 until Task A8 wires the real enum and the streaming writers through.
StreamingCorrelatorTest gained an auditLevels case for that same CLI
contract, but it is not called from main: the report has no "audit" field
until A8 lands, so calling it now would fail on an unimplemented assertion
rather than a missing feature. A8's dispatch carries the same instruction
to re-enable it.
…ysis facade

analysis() and correlatePartial handed the Collector's live ArrayLists straight
into the Analysis/PartialAnalysis constructors, so analysis.records() and
analysis.matches() were mutable and mutating records() also mutated the
JsonObject instances Match aliases. The deleted retained path always wrapped
both in List.copyOf; restore that at both call sites.
The compatibility writer was the last consumer of a retained observation and sample
per match, and it re-derived a canonical stack key for each one. It now takes the
matched intervals as primitives and reads their frames and threads out of the
correlation's dictionaries, so the same events are written from a table of distinct
stacks rather than from a copy per sample.
…nd tie-break

SyntheticJfrSource.of(Analysis) had silently traded four checked IOExceptions
(missing correlationId/frames/monotonicTimeNanos/startTime, and the old
validateMatch's null/relational check) for unchecked NPEs/ArithmeticExceptions.
Restored all of them, plus checked long-narrowing for the from/to/duration/epoch
fields that now have to fit primitive longs.

Extended syntheticFromColumns to compare actual per-event synthetic JFR timestamps
between the column-backed and retained plans (previously only counts and totals were
checked, so a dropped or sign-flipped epoch offset would have passed silently), and
added syntheticOrderTiesBreakOnCookie to pin the merge sort's cookie tie-break with
two matched intervals sharing identical fromNanos/toNanos submitted in descending
cookie order. Verified both new assertions fail on a deliberately reintroduced bug
before confirming they pass on the fix.
One writer now serves both the streamed CLI path and the retained library view, so
the two cannot drift, and --audit selects how much per-row output to produce.
jonoffcpu-classified-records.jsonl and jonoffcpu-matches.jsonl are the most expensive
outputs and the least often read in full; a fifteen-thousand-row run once wrote 71 MB
of them.

The digest recheck (verifyUnchanged) is threaded into the writer as a callback invoked
after every other artifact -- including the audit re-read for the streamed path -- but
before the completion marker, so a file changed under the correlator during the second
pass is still caught before the directory is promoted to complete.
Two million observations over a couple of thousand distinct stacks now correlate in a
one-gigabyte heap with peak retained bytes asserted below a fixed bound, which is the
property the columnar engine exists to provide.
The expansion cap used to reject an analysis once the collapsed stacks were already
on disk, which threw away work that was correct and left a four-minute benchmark to
be run again. The event count for a quantum is the sum of per-stack floors and is
known before the first event, so the quantum is raised until it fits and the raise is
recorded in the report and in the recording's metadata event. A synthetic JFR renders
totals the report states exactly, so a coarser quantum costs granularity, not time.
…ighting

A capture that does not fit can still answer the question it was taken to answer. Each
interval is kept with probability q, decided by hashing its cookie so the choice is
deterministic, order-independent and reproducible, and the duration it contributes is
scaled by the exact reciprocal of the realised probability. Because the cookie is the
join key, an observation and its sample are dropped together and every count in the
report keeps describing one coherent subsample.

This changes what the collapsed weights are: observed durations become estimates. The
report carries sourceThinning with the estimator and the exact realised probability,
the collapsed stacks carry a root label so a rendered flame graph cannot be quoted
without its caveat, and q = 1 stays the default whenever the input fits.
…ale()'s stack fan-out

ScaleFixture.recording's row-to-stack skew, added for the thinning-accuracy fixture's top-5
comparison, had replaced the shared uniform round robin scale() also depends on, silently
collapsing scale()'s ~2000-stack fan-out coverage down to ~21 distinct stacks while its only
check (stackCount() <= 4000) kept passing vacuously.

Splits the skew into its own recordingSkewed() helper used only by thinningAccuracy; scale()
keeps the original uniform round robin. Adds a floor assertion to scale() so this class of
regression is caught going forward; the floor is not an exact count because the JVM's default
JFR stack-capture depth (64 frames) plus JIT/inlining-dependent truncation shape make the
fixture's real fan-out (measured at 165 and 171 distinct stacks across two runs) inherently
non-deterministic, never the full 2000 requested.
Retention per record is a constant of the columnar design and the file sizes are
known, so the correlator can choose its thinning factor and audit level at the start
and report the decision, instead of discovering at ninety percent of heap that the
input does not fit. The budget failure is now a type rather than a message, so the
ladder can retry a limit while a semantic failure still aborts.
…done

A capture that does not fit used to produce nothing, even when the useful
output was already on disk: one run wrote a collapsed-stack file, a usable
off-CPU flame graph, then failed on the synthetic JFR and left nothing
behind. The cost of that refusal was not a wrong number, it was work that
had to be redone from scratch.

Integrity failures still abort, because they mean the two files do not
describe the same capture. Volume failures now walk a ladder -- drop the
audit outputs, thin the source and reweight, narrow the window -- and every
step taken is named in the report's degradation object. Thinning keeps the
ordinary names and exits zero because a stated estimator over the window
that was asked for still answers the question; narrowing takes the
INCOMPLETE names and exits two because it does not.

--max-rows now defaults to a hundred million and --max-retained-bytes to a
share of the heap it protects, since neither was a meaningful guard while
retention scaled with input size.

Two fixes to CorrelationEngine were necessary for the ladder to actually
converge rather than loop or refuse unnecessarily:

- A narrow-window cut now excludes the triggering observation by its own
  start, not its end. Cutting at the end left that observation itself
  admitted on retry, so a restart retraced the identical column growth up
  to the identical watermark and reproduced the identical failure -- an
  infinite loop, confirmed by a hung run before this fix landed.
- The JFR-pass watermark can now also hand back a cut (via the failing
  sample's own source slot) instead of always throwing with none. Without
  this, any capture whose JFR samples contribute materially to retention
  could never actually finish narrowing: the watermark that first exceeds
  budget during the source pass necessarily measures less than the eventual
  source-plus-sample total, so a cut taken only from a source-side
  exceedance is mathematically guaranteed to still be over budget once
  samples are added, and previously that second exceedance had no cut to
  give the ladder.

Also sizes the sample-side columns off the JFR file's own size rather than
the capture file's, so a JFR much smaller than its capture does not
preallocate as though it were exactly as dense.
…cut,

cover the watermark-driven thin-source rung, and clean up refusal/report text

- drop-audit-outputs changed nothing the engine's own retention accounting
  sees (audit output comes from a separate post-hoc re-read), so returning
  after one level meant the caller's next attempt was a guaranteed-identical
  repeat failure on a full pass. Now drops every remaining audit level in
  one advance() call, still recording each as its own step, before falling
  through to a step that can actually change the retry's outcome.
- An orphan JFR sample (no source row at all, distinct from a source row
  DROPPED by thinning or narrowing) used to rethrow with no cut point when
  it happened to land on an exceeded watermark, costing a rung it did not
  need to. Falls back to the last kept source slot's start instead, since
  the source pass is always complete by the time samples stream.
- Added a ladder-reactive-thinned fixture that forces Degradation's retry
  loop, not just its constructor, to pick a thinning rung: the constructor
  already picked a sufficient rung for the tuned budget on attempt 1, so the
  existing assertion never touched advance()'s thin-source branch. Confirmed
  the gap first (a debug print showed attempts=1 against the old fixture).
- Degradation.refusal() no longer tells a truncate run to narrow --from-ns/
  --to-ns by hand when it has already narrowed as far as the watermarks
  allow, and no longer claims degrade would thin the source when the caller
  already requested thinning explicitly.
- Degradation.report() gains a top-level narrowedToNanos, symmetric with
  peakRetainedBytes, so a consumer does not have to scan stepsApplied and
  take the minimum.
…dder

--max-retained-bytes now charges what is retained rather than what was decoded, which
makes it a number a caller can plan a capture against: about 80 bytes per recorded
interval plus one copy of each distinct Java stack, measured at ~281 MiB peak for the
2,000,000-observation scale fixture. The two kernel-proof tools that read
jonoffcpu-classified-records.jsonl now ask for it explicitly, since --audit defaults to
matches on the CLI (the library default stays full). Also documents --thinning,
--on-limit, the five-step degradation ladder and its report object, the
thinning-vs-narrowing output-naming/exit-status distinction, and the signed 64-bit
nanosecond bound on --from-ns/--to-ns/--max-handler-delay-ns.
…ixture

Task E1 measured 212,831,820 bytes (203 MiB) of peak retained bytes
correlating the real Pulsar broker capture at -Xmx2g, -Xmx512m and
-Xmx14g alike (~190 bytes/recorded interval), about 2.4x the documented
80-byte column-only figure from the synthetic scale fixture. Note the
real-world number next to it so the 80-byte figure isn't read as the
full per-interval budget.
`--audit full` combined with any row-dropping degradation crashed with
ArrayIndexOutOfBoundsException. AuditPass was built on a 1:1 file-row to
column-slot assumption; thinning and window narrowing later made the columns
hold only kept rows, so the mapping slipped at the first dropped row and then
ran off the end of the columns. It died after the collapsed stacks and the
synthetic JFR were already written, with a raw JVM stack trace rather than the
IOException run(args) documents.

Teach the audit pass the kept-row predicate instead of forcing the audit level
down: an audit describing the kept subsample is what --audit full should mean
under thinning. CorrelationResult carries the narrow cut alongside the thinning
it already held and exposes keepsSource/keepsSample, whose static overloads
CorrelationEngine's own two admission points now call — the streaming pass and
the re-read are one function, so they cannot drift apart again. The require()
backstops are kept and are now true invariants of the new mapping.

StreamingCorrelatorTest.auditUnderDegradation covers both mechanisms, asserting
that the audit rows are the rows actually kept: counts against the report's
classification counters, and the matched cookie set against the matches file.
It is called after scale(), not next to auditLevels(), because scale()'s
distinct-stack floor collapses when another large fixture warms ScaleFixture's
recursion ahead of it.

Also: remove code orphaned by later tasks (CaptureInput.frames, two write-only
CaptureInput fields, the unsigned comparison helpers nothing calls and their
assertions); drop the stale merge-window mention from Budget's javadoc; stop
--on-limit truncate claiming the window is fully narrowed when no watermark ever
offered a cut point; document the --estimate-population/thinning asymmetry, what
an audit of a degraded run describes, and the scale fixture's retention as the
bound the test actually asserts rather than a drifting measurement.
@lhotari
lhotari merged commit 325188d into main Sep 22, 2026
5 checks passed
@lhotari
lhotari deleted the correlator-streaming-join branch September 22, 2026 21:01
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant