Conversation
Preserve typed OTLP attributes and UTC nanoseconds in format-3 telemetry, use native extraction and timestamp pushdown, bound keyed dependency traversal, and verify Parquet with a reusable native engine. Reject prior formats without rewriting them and keep SQLite control state separate.
vishr
left a comment
There was a problem hiding this comment.
Review: DuckDB 2 adoption and why the benchmark did not move
The DuckDB 2 integration is sound where it is used: timestamp predicates stay on native TIMESTAMPTZ_NS columns and prune, the views keep union_by_name/Hive inference off, VARIANT is shredded and the hot messaging projection reaches the scan, and async I/O is correctly left at its memory-governed defaults. The benchmark is flat because the engine is not the bottleneck. Details and inline findings below.
What I measured
I replayed the exact dashboard SQL and bound arguments the benchmark issues (dumped from each source tree) against the retained benchmark datasets, inside the same 4-CPU/6 GiB Docker VM, with no ingest running.
1. Every dashboard statement is fast in isolation. The slowest statement in either build is about 0.3 s (baseline endpoints query 288 ms; PR endpoints query 169 ms, recent-trace query 173 ms). Under the benchmark, /api/observability/performance had p50 8.1 s. The roughly 40x gap is queueing on a saturated CPU and the 4-connection pool, not statement speed.
2. The 20 reads/s gate is not reachable on this host with either engine. With 4 closed-loop workers, no ingest and all 4 CPUs, the full dashboard mix peaks at about 8.2 requests/s on the baseline and 9.5 requests/s on this PR (+16%). The mix costs about 1.8 CPU-seconds per five requests on the PR (/performance ~0.63, /trace ~0.73, /logs ~0.42). At 20 requests/s that is about 7 cores of demand on a 4-vCPU VM that also runs telemetrygen and ingest. All three builds completed about 7.4/s. In the first mixed phase Fanout used 3.3-3.7 cores, versus about 0.2 cores for write-only ingest at a similar rate, so reads are what saturate the box.
3. Engine-only comparison on identical data and SQL is roughly flat (baseline format-2 dataset and baseline SQL, threads=4, median of 5):
| Statement | DuckDB 1.5.5 | DuckDB 2 alpha |
|---|---|---|
/performance endpoints (rollup + raw boundary) |
288 ms | 225 ms |
/trace recent-trace selection |
224 ms | 244 ms |
/logs severity histogram |
80 ms | 78 ms |
/logs entries |
8 ms | 14 ms |
| Rollup-table reads (overview, topology, points) | ~1 ms | ~2-3 ms |
The PR's faster endpoints query on its own dataset (169 ms) comes mostly from the format-3 layout and a smaller dataset, not the engine. On a generic suite over the same data, DuckDB 2 is clearly faster only on some shapes: wide ORDER BY ... LIMIT 100 (297 ms to 18 ms), grouped quantiles (-36%), parent/child self-join (-22%) and windows (-11%). High-cardinality GROUP BY, COUNT(DISTINCT), sorting, string search and regexp are flat. The large DuckDB 2 wins (async I/O for remote Parquet, recursive CTE USING KEY, VARIANT shredding) are not on the dashboard hot path. Parquet is written by parquet-go, so DuckDB's writer changes do not apply either.
4. Per-request work scales with all retained data. The benchmark's 6h/12h/24h windows exceed the dataset's age (about 6 minutes), so every request scans about 5M rows. The inline comments on performance.go, trace.go and logs.go cover the three statements that account for nearly all read CPU.
Findings outside the diff
- Medium:
.github/workflows/site.yml:80will fail.go run ./cmd/fanout-docgen --checkruns withoutscripts/with-duckdb.sh, but docgen now linksinternal/duckdbthroughinternal/query. In a scratch copy it fails withinternal/duckdb/native.go:6:10: fatal error: 'duckdb.h' file not found. The justfile recipe was wrapped; this workflow step was not. - Low:
ensureLimit(internal/query/sql.go:387) is now dead.projectSQLResultsowns theLIMIT; only tests call it. - Context for the inline
performance.gocomment:rollupPublicationSafetyLag = 5 * time.Minute(internal/query/duck.go:136). At the end of the final run, the endpoint rollup watermark was 4.6 minutes behind the newest raw span.
Benchmark methodology
The current harness cannot show a read-path improvement on a 4-vCPU same-host VM: the read gate is infeasible for every build, and sweep-phase ingest volume (which differs per build) sets the dataset size for the sustained phase. A meaningful comparison needs generators on a separate host (or a read rate scaled to the cores), windows that match the dataset age, and at least three repetitions. Many of the reported "query errors" are 15-second client cancellations logged as server 500s; see the inline comment on mapQueryError.
| AND (s.start_time < ? OR s.start_time >= ?) | ||
| WHERE s.start_time >= ?::TIMESTAMP_NS::TIMESTAMPTZ_NS | ||
| AND s.start_time < ?::TIMESTAMP_NS::TIMESTAMPTZ_NS | ||
| AND (s.start_time < ?::TIMESTAMP_NS::TIMESTAMPTZ_NS OR s.start_time >= ?::TIMESTAMP_NS::TIMESTAMPTZ_NS) |
There was a problem hiding this comment.
Performance: this raw boundary scan is most of the read CPU, and it stays large in steady state.
The interior bound is capped at the endpoint rollup watermark, which trails ingest by rollupPublicationSafetyLag (5 minutes, internal/query/duck.go:136). Every /performance request therefore scans at least the newest five minutes of raw spans and evaluates 17 COUNT(*) FILTER buckets per row.
At the end of the final run the watermark was 4.6 minutes behind the newest span, so this covered most of the dataset. The endpoints statement cost about 0.63 CPU-seconds per request (5.4M spans), versus about 0.01 for the five rollup-table statements combined. The cost is proportional to ingest rate times the lag, not to the requested window, so a busy production instance pays it on every request.
Options: shorten the lag for the endpoint cache specifically, cache the boundary aggregate per minute between rollup passes, or serve the partial tail from a cheaper histogram.
There was a problem hiding this comment.
Agreed; this remains open performance work. I kept the five-minute publication safety lag rather than changing correctness/freshness policy to improve the benchmark. The bounded endpoint response does not bound the raw-tail scan. A follow-up should materialize/cache per-minute endpoint histograms for the unrolled region and merge them with published aggregates, retaining exact partial-minute boundaries, namespace/service isolation, and invalidation for new or late-arriving batches. It needs parity checks against the current raw query and repeated measurements at a fixed ingest rate. The current revision does not remove the raw boundary scan.
There was a problem hiding this comment.
Agreed. Keeping the five-minute publication safety lag is the right call; the fix belongs in caching the unrolled region, not in loosening freshness. The plan in your reply (per-minute histograms for the tail, exact partial-minute boundaries, invalidation for late batches, and parity checks against the raw query) covers what I'd want. Please track it in an issue. Leaving this open.
| FROM spans | ||
| WHERE start_time >= ? AND start_time < ? AND (? = '' OR namespace = ?) AND trace_id <> '' AND (? = '' OR service = ?) | ||
| WHERE start_time >= ?::TIMESTAMP_NS::TIMESTAMPTZ_NS AND start_time < ?::TIMESTAMP_NS::TIMESTAMPTZ_NS AND (? = '' OR namespace = ?) AND trace_id <> '' AND (? = '' OR service = ?) | ||
| GROUP BY trace_id |
There was a problem hiding this comment.
Performance: choosing one trace aggregates every trace in the window.
With no trace_id, this hash-aggregates all spans in the 12h window by trace_id (about 2.7M groups on the benchmark dataset) to return one row. It costs about 0.7 CPU-seconds per /trace request (173 ms with 4 threads), about as much as the whole /performance request.
I tried a root-span top-N (parent_span_id = '', ordered by error then duration_ms, LIMIT 1): 103 ms, but it only saves about 20% CPU because it still scans every row. A real fix is a precomputed candidate, such as a small "notable traces" table maintained by the rollup pass, or a short default window for the "pick a trace" case.
There was a problem hiding this comment.
Agreed; leaving this item open. LIMIT 1 bounds the answer, not the high-cardinality aggregation needed to choose it. A root-span shortcut or shorter implicit window would change the current selection contract (error priority, elapsed bounds across selected spans, service/namespace scope, and incomplete/late traces). The follow-up should maintain a notable-trace candidate index with late-span updates and a bounded current tail, then verify its choice against the existing query across scoped and partial windows. No candidate precomputation or default-window change is included in this revision.
There was a problem hiding this comment.
Agreed. A root-span shortcut does change the selection contract, which is why I'd also prefer a maintained candidate over it. The notable-trace index with late-span updates and parity checks against the current query sounds right for a follow-up issue. Leaving this open.
|
|
||
| var logBucketsQueryTemplate = ` | ||
| SELECT time_bucket(INTERVAL '%s', time) AS point_time, | ||
| SELECT time_bucket(INTERVAL '%s', time::TIMESTAMP_NS) AS point_time, |
There was a problem hiding this comment.
Performance: the log histogram re-scans the full window on every request.
For the benchmark's 24h window, this reads time and severity for every log row (4.96M on the final dataset): about 0.39 CPU-seconds per /logs request, against 0.03 for the entries query beside it. Per-minute counts by severity are a natural rollup, like service_rollup.log_count, with only the unrolled tail read from Parquet.
There was a problem hiding this comment.
Agreed; this remains open performance work. Corrected the misleading comments in 9d6a7c4: LIMIT/GROUP BY bound the response size, while histogram work still grows with matching retained rows. A severity rollup needs (minute, namespace, service, severity) dimensions, published-interior plus exact raw boundaries/tail, and late-arrival handling. Arbitrary searches match redacted body text, so those counts must still use the matching raw predicate; a severity-only rollup cannot answer them. The follow-up should prove entries, buckets, and matched totals remain consistent across these cases. No persistent log histogram rollup is implemented here.
There was a problem hiding this comment.
Agreed, and the corrected comment is accurate now. Good point that free-text searches match redacted bodies, so only the unfiltered and severity-filtered cases can come from a rollup. Follow-up issue. Leaving this open.
| // writeTelemetryColumns streams variant events straight into shredded columns. | ||
| // The generic map encoder builds a second object tree for every row; keeping | ||
| // the canonical OTel maps as the event source bounds the ingest working set. | ||
| func writeTelemetryColumns[T any](writer *parquet.GenericWriter[T], rows []T) error { |
There was a problem hiding this comment.
Performance (not profiled; worth measuring before merge): ingest cost per row.
Peak mixed throughput fell from 145.9k to 117.1k rows/s (-20%) and cooldown RSS rose from 0.87 to 1.54 GiB in the PR's benchmark. That is one trial on a CPU-saturated host, so it is not conclusive. Plausible contributors are the per-row reflection plus variant encoding here, and ValueBytes charging about 64-80 B per attribute key while re-walking the shared resource map for every row. A CPU and heap pprof of a write-only mixed-p8 phase against #270 would settle whether this path or the benchmark noise explains it.
There was a problem hiding this comment.
Measured in a separate ingestion-only experiment: the existing durable shared-transport fixture, 16 concurrent clients, 1,000 rows/export, equal spans/logs/metrics, no dashboard reads; three alternating baseline/current repetitions on the same native Linux ARM64 4-CPU/6-GiB container. CPU profiles cover each ten-second load interval; heap profiles are taken before cleanup/verification. This is not the harness's mixed-p8 workload and does not establish its cooldown RSS cause.
| Median | #270 | Typed DuckDB 2 branch |
|---|---|---|
| Durable acknowledged rows/s | 484,312 | 411,796 |
| Go allocated bytes/row | 3,826 | 4,123 |
| Client + server CPU cores | 1.316 | 1.478 |
| Export p95 | 36.54ms | 42.09ms |
All six runs passed storage verification with zero export errors. A roughly 15% throughput regression reproduces on this fixture, so it cannot be dismissed as dashboard contention alone. The first current profile attributes about 2% cumulative CPU to ValueBytes; its per-key charge is conservative accounting, not an allocation of that size. Typed Parquet writing, encoding, copies, and GC contribute more. Added optional measurement-window CPU/post-load heap profiling to this existing experiment in 9d6a7c4, with profiles separated by subtest. Updated the PR's results and limitations. The writer regression and original retained-memory result remain open performance concerns; I have not changed admission accounting or claimed a fix.
There was a problem hiding this comment.
Thanks for measuring this properly: three alternating runs with profiles settle it better than the single harness trial did. I have not reproduced it independently.
That makes the merge decision explicit: about 15% lower durable ingest throughput, about 12% more CPU and 8% more allocation per row, in exchange for typed attributes. I'd attach these profiles to a follow-up issue for the writer path, and re-measure cooldown RSS with repeated runs. Leaving this open as the known regression.
|
Follow-up 9d6a7c4 fixes five inline correctness/build/error-classification findings, plus the site docgen workflow and dead SQL-limit helper from the review summary. The pre-commit lint job also now uses the pinned native wrapper. Added regression tests and reusable opt-in ingestion profiling. Local |
What changed
Fanout now embeds DuckDB
v2.0.0-alpha43763and stores OTLP attributes as typed, shredded Parquet VARIANT values. Integers, booleans, bytes, nested values, and nanosecond UTC instants survive ingestion and compaction instead of being flattened into JSON strings.SELECT version().Breaking contracts
Telemetry batch format 3 and its physical schema are the only accepted format. Startup rejects unsupported metadata before cleanup or schema rewriting and preserves those files. There is no legacy reader, automatic conversion, route alias, or format fallback. Start with an empty telemetry directory and retain prior telemetry for its original reader. SQLite control state remains on modernc/database/sql, sqlc, and Goose.
SQL uses typed
attributesandresourcecolumns instead ofattributes_jsonandresource_json;attr()returns a typed VARIANT and callers cast scalars explicitly. General VARIANT extraction across view projections remains a limitation of this pinned preview; the hot messaging projection has a verified scan-level workaround.Verification
just checkjust test-race(full Go suite; existing documented checkptr workaround for the CGO driver)Additional local validation: native Linux ARM64 query/intelligence/observability and telemetry/storage tests; runtime engine version; typed-Parquet interoperability and compaction; SQL boundary restrictions; nanosecond/DST predicates and executed row-group pruning; executed shredded messaging projections; detector/log-helper queries against production Parquet; late-page corruption, span ordering/index parity, cancellation/OOM quarantine safety; generated docs (13 pages); UI tests, dependency audits, and site build. Native macOS ARM64 was tested locally; the new CI matrix also exercises AMD64.
Review follow-up
Benchmark results and limitations
Matched standard harness phases on Linux ARM64, 4 CPUs / 6 GiB, same-host telemetrygen
v0.157.0, 20 requested reads/second, 15-second generator sweeps and 60-second sustained mixed load. One trial per build; accepted dataset sizes vary. The dashboard windows (6h/12h/24h) exceeded the roughly six-minute dataset age, so these are whole-dataset saturation measurements. They are not representative window-pruning or capacity comparisons. The final source tree SHA256 was797cb518c8310bd54e544575caa7df46268801f1a99b68c219127b474bcc8688and remained unchanged during the run. Baseline is the #270 merge,30ba9b9ba8a4217cc73637760d7acdcd5f89f15a.All builds passed storage verification with zero server-dropped rows, restarts, or generator failures. All builds failed the sustained 20 reads/second quality gate. The final build restores completed reads near baseline but has lower mixed peak throughput and higher retained memory in this single run. These results do not establish a general capacity improvement. Most logged HTTP 500s in these original runs were 15-second client cancellations; the new HTTP 499 mapping corrects server-side classification but does not make those requests successful or change the historical counts. A capacity comparison needs separate generators or a read rate that fits available CPU, windows appropriate to dataset age, identical fixtures, and at least three repetitions.
A separate same-batch full-validation comparison improved from 5.88s to 0.34s (single run each, both including CPU profiling overhead). Rechecking the original initial-DuckDB-2 retained dataset with the final verifier passed in 18.3s, including Docker startup. These verification checks are distinct from the ingestion/read benchmark.
A separate ingestion-only check used the existing durable shared-transport fixture (16 concurrent clients, 1,000 rows/export, equal spans/logs/metrics, no dashboard queries), with three alternating baseline/current repetitions on the same 4-CPU/6-GiB Linux ARM64 container. CPU profiles cover the ten-second measurement intervals; Go allocations and CPU include the in-process clients and server. This is not the harness's
mixed-p8workload and does not measure its cooldown RSS.All six runs passed durable storage verification with zero export errors. The approximately 15% throughput regression reproduces on this fixture; it cannot be dismissed as read-load noise. In the first current CPU profile,
ValueBytesaccounts for approximately 2% cumulatively, while typed Parquet writing, encoding, copying, and GC contribute more. The conservative per-key byte charge is accounting, not a corresponding allocation. The writer regression remains open; the retained-memory result from the original mixed benchmark is still a single-trial observation.Relevant upstream references: DuckDB 2 preview, async I/O memory policy, and pinned optimizer projection behavior.