Skip to content

parquet: json_parquet_encode allocates a third as often, reuses writers, caches by metadata - #5

Open
almostintuitive wants to merge 20 commits into
perf/json-parquet-encodefrom
perf/json-parquet-encode-allocs
Open

almostintuitive wants to merge 20 commits into
perf/json-parquet-encodefrom
perf/json-parquet-encode-allocs

Conversation

@almostintuitive

@almostintuitive almostintuitive commented Sep 29, 2026 •

Copy link
Copy Markdown

Summary

Stacked on #3. Follow-ups to json_parquet_encode, one commit each, each checked against the replaced pipeline by the parity tests and the fuzz target. They're followed by fixes for everything an adversarial review of the first version found.

  1. Drop the write buffer and per-row key strings.
  2. Parse and coerce without allocating. Keys and fields are read as raw JSON. A string or number is coerced by the same strconv call the Bloblang coercion makes; bools and nested values still go through the original helpers.
  3. Keep each partition's rows in its own arena. Each partition's values and bytes come from chunks that start at one row and 64 B and double up to 256 rows and 16 KiB. They're freed when the partition's file is written.
  4. columns[].cache_by (opt-in). An interpolated column is evaluated once per batch for each distinct combination of the listed metadata fields.
  5. Reuse writers through a pool, reset for each file.
  6. Fewer, larger files. The benchmark now covers one file, few, many, and one per message. The docs explain what each file costs, and a json_parquet_encode_files counter makes the file count visible in production.

Review fixes

  • Long strings were corrupted. A text value over 4 KiB was "copied" with []byte(b) of a []byte, which doesn't copy. The next message parsed would overwrite it, including with bytes from other partitions. It now always copies. A new parity test and new fuzz seeds send long strings among other messages; both fail without the fix.
  • cache: true could key the cache on too few fields. Bloblang doesn't report everything a value reads: a call with a non-literal argument reports only its arguments, and a lambda only its receiver. The option is now cache_by: [fields], and the cache is keyed on exactly those fields. What Bloblang does report is still checked, so reading a message field, the whole metadata, or an unlisted field is rejected. What it can't find is listed in the docs. The public XUnwrapper is gone; the expression is parsed through the manager's Bloblang environment, as group_by_value does.
  • An optional interpolated column was always NULL (predates this branch). It now uses the column's definition level.
  • Reused writers. The actual cause is that parquet-go 0.29.0's byte-array column buffer keeps its scratch buffer across Reset. On a reset writer, a page with a single empty value writes empty min/max, which a new writer leaves out. A file now gets a new writer only when a text column has exactly one non-null value and it's empty. As a precaution it also does when the file could span pages and contains an empty string; the tests couldn't show that second case is needed. A new differential test covers every encoding and codec, partitioned and not. It encodes the same batches with pooled writers and with a new writer for every file, and requires identical bytes and full reuse on clean batches. Without the check, 15 of its 36 configurations fail. Writers are reset to nil before going back to the pool, so they don't keep their last file reachable.
  • Tests:
    • The count() test failed under -count=2, because Bloblang's counters are process-global; it now uses a new counter each run.
    • The escaped-field-name case replaced a name with itself; it now really escapes names.
    • The parity harness and the benchmark enable the cache through one helper that finds the interpolated column, instead of assuming it's the first.
  • Memory:
    • Arenas start at one row, so tiny partitions no longer hold several times their rows.
    • A rejected row no longer creates and discards a group and arena chunks.
    • Files are written into one buffer per batch and copied out at their size, instead of each keeping what doubling left unused.
    • Strings the processor built itself are no longer copied.
  • Speed: column names that need no unescaping are looked up as they appear, without a scan. Parsing a message takes about 14% less time.

Results

10 MiB batch, median over repeated runs on a 16-vCPU AMD EPYC. The PR's pipeline caches the interpolated column, as the docs example does.

base branch this PR
26 files: time / allocated / allocations / peak heap 426 ms / 104 MB / 2.30M / 75 MiB 312 ms / 69 MB / 0.76M / 59 MiB
710 files 661 ms / 202 MB / 2.61M / 101 MiB 436 ms / 93 MB / 0.91M / 80 MiB
one file — 319 ms / 95–102 MB / 0.72M / 81–89 MiB
one file per message (35,966 files) — 5.3 s / 1,177 MB / 8.49M / 223 MiB (1,484 MB / 277 MiB before the review fixes)

The legacy pipeline on the same batches takes about 2 s, allocates 485–605 MB and peaks at 184–231 MiB. A batch of 10,000 rejected rows allocates 6.9 MB (10.3 MB before the review fixes).

Things to know

  • cache_by is a contract. The value must depend on nothing but the listed fields. Functions that read the message without naming a field (content()) or that change between calls (now(), count()) can't be detected.
  • Evaluation order. A row's partition path is built before its values are coerced. When both would reject a row, only the logged error changes.
  • Upstream fix. Fixing parquet-go's Reset upstream would remove the need for new writers entirely. That isn't filed yet.

Test plan

  • go test -count=2 ./internal/impl/parquet/, and go test ./internal/bloblang/... ./public/service/ ./internal/impl/kafka/
  • Mutation checks, each caught by the tests:
    • parsing integer strings in base 10;
    • skipping string unescaping;
    • an absent field taking the previous message's value;
    • reusing a writer regardless of empty strings;
    • dropping the one-value rule;
    • looking up names without unescaping them;
    • the long-string copy.
  • Fuzz after every commit (30–180 s), plus 5 minutes at the end: no divergence
  • golangci-lint v2.13.2 on all touched packages: 0 issues
  • Docs regenerated with bento_docs_gen

🤖 Generated with Claude Code

almostintuitive and others added 19 commits September 29, 2026 13:53
…trings

The writer buffered its output in 32 KiB before writing it to a bytes.Buffer,
once per file; it now writes straight to the buffer, with the same bytes. The
partition path is appended to a reused byte slice, and a string is only made
of it for a new partition.

10 MiB batch across 710 partitions: 200 MB -> 175 MB and 2.61M -> 2.44M
allocations per batch, peak heap 115 -> 101 MiB, 680 -> 636 ms.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Every key and every field the schema reads was made a string, and every field
boxed into an interface, only to be looked up, parsed and dropped. Keys and
fields are now read as raw JSON; a key is looked up by its bytes, and a field
is copied into a per-message buffer as a string's decoded text or a number's
literal. A string or number is coerced by the very strconv call the Bloblang
coercion makes of it (json.Number's Int64 and Float64 parse in base 10, a
string's integer in any base Go literals use); a bool or nested value is still
decoded into the Go value Bento's decoding gives and handed to the coercion.
A string is decoded as jsontext.Token.String decodes it: its own bytes when it
has no escape and is valid UTF-8, AppendUnquote's mangling otherwise.

TestJSONParquetParityOptionalPartitionField covers an optional field in the
partition path that is present, absent and null in turn, which a mutant
reading an absent field as the previous message's otherwise got past.

10 MiB batch across 710 partitions: 2.44M -> 1.68M allocations and 175 ->
161 MB per batch, peak heap 101 -> 94 MiB, 636 -> 560 ms.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The encode benchmark only split its 10 MiB batch into hundreds of files,
where a file's own costs (its writer, its footer) dominate. It now also
partitions the batch as the archiver example does, by exchange and day, into
a couple of dozen files, and reports the file count.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Each kept row was a slice of its own, and each of its UTF8 values another.
A partition's rows now take their values and bytes from chunks that double
from a few rows up to 256 rows and 16 KiB, dropped with the partition once
its file is written, so the rows still shrink as the files are encoded (one
arena for the whole batch held every row until the last file, and peaked
10 MiB higher). A new partition is only kept once one of its rows is, and a
dropped row gives back what it took.

A row's partition path is now formatted before its values are coerced. Which
of the two rejects a row that both would only changes the error logged.

10 MiB batch (per batch, against the previous commit):
- 26 files:  1.40M -> 1.17M allocations, 88.0 -> 87.3 MB, peak heap 69 -> 68 MiB,
             384 -> 366 ms
- 710 files: 1.68M -> 1.46M allocations, 160 -> 168 MB, peak heap 93 -> 98 MiB,
             528 -> 526 ms; the last chunk of many small partitions is part-empty

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A column set from an interpolation, such as the exchange read from the
Kafka topic, ran its Bloblang for every message: about a quarter of the
processor's allocations. With cache: true it is evaluated once per batch for
each distinct combination of the metadata fields it reads, found from the
query's targets. A value that reads a message field or the whole of the
metadata is rejected when the config is parsed. A function that reads the
message without naming a field (content()) or that differs from call to call
(now(), count()) has no targets and cannot be told apart, as the field's
description says. A message whose field holds anything but a string is
evaluated on its own.

field.Expression reports its query targets, and service.InterpolatedString
gains the XUnwrapper other public types have, to reach them.

The encode benchmark caches the column, as the archiver example now does.
10 MiB batch, against the previous commit:
- 26 files:  1.17M -> 0.77M allocations, 87 -> 78 MB, peak heap 68 -> 63 MiB,
             366 -> 333 ms
- 710 files: 1.46M -> 1.06M allocations, 168 -> 156 MB, peak heap 98 -> 96 MiB,
             526 -> 495 ms

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Making a writer for every file was a third of what a batch of many small
files allocated. Writers are now kept in a pool and reset for the next file.

A reset writer writes the bytes a new one would, save for one case, found
by the parity tests: a page whose values are all empty strings. A new writer
leaves its bounds out, since they point into a column buffer that holds
nothing, and a reset one writes them as empty, since its buffer held the last
file's values (parquet-go 0.29.0). A file with an empty UTF8 value therefore
still takes a new writer; without that, five parity tests fail.

10 MiB batch, against the previous commit:
- 710 files: 156 -> 119 MB and 1.06M -> 0.95M allocations per batch,
             495 -> 438 ms, peak heap unchanged (96 MiB). A quarter of the
             benchmark's symbols are "", so most of its files still take a
             new writer.
- 26 files:  unchanged within noise.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Every file costs something of its own: a footer, pages begun for each column,
and, with an empty string, a new writer. On the encode benchmark the same
10 MiB batch takes 438 ms and peaks at 96 MiB as 710 files, against 331 ms and
62 MiB as 26. The docs now say so and advise a coarse partition path, and the
files written are counted in json_parquet_encode_files, so that how many a
batch makes can be seen in production. A test covers both counters.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A UTF8 value over 4 KiB was "copied" as []byte(b), which, with b already a
[]byte, is b itself: a view into the buffer the next message is parsed into.
The value was then overwritten by the fields of later messages, of any
partition, before its file was written. Found by review; the parity tests only
sent a long string in a batch of its own.

TestJSONParquetParityLongStringsAmongOthers sends strings of every size
around the arena's among other messages, partitions and dropped rows, and the
fuzz target gains seeds of the same shape; both fail without the fix.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…lds listed

cache: true keyed the cache by the metadata fields Bloblang reported the value
reads, and Bloblang does not report them all: a function or method called with
a non-literal argument reports only its arguments' reads, and a lambda only
its receiver's. So @kafka_topic.split(@sep).index(2) was cached by sep alone,
and messages of different topics shared a value, and a partition file (found
by review).

The option is now cache_by, the metadata fields the value depends on, and the
cache is keyed by exactly those. What Bloblang does report is still checked
against them: a value found to read a message field, the whole of the
metadata or a field not listed is rejected, which catches the mistake above
(split(@sep) with cache_by [kafka_topic] reads sep). What it cannot find is
listed in the field's description, which no longer claims more.

The expression is parsed through the manager's Bloblang environment, as
group_by_value and dedupe do, so service.InterpolatedString no longer gains an
XUnwrapper. field.Expression.QueryTargets switches on the resolver type
instead of asserting an anonymous interface.

Tests: TestJSONParquetEncodeCachedByEveryField shows a value with a
non-literal argument is right when every field is listed; rejections cover
an unlisted field and an unlisted non-literal argument. The count() test
takes a new counter each run, since Bloblang's counters live for the process
and it failed under -count=2. The parity harness and the benchmark cache the
interpolated column through one helper that finds it, rather than assuming
it is the first.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… value

A column set from an interpolation was always written at definition level 0,
which in an optional column is NULL: every row read NULL where the mapping it
replaces wrote the value (found by review; the bug predates this branch).
It now takes the column's definition level, as a field's value does.

TestJSONParquetParityOptionalInterpolatedColumn makes the archive schema's
interpolated column optional, cached and not, and fails without the fix.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ould differ

Any file with an empty UTF8 value took a new writer, on a wrong account of
why a reset one differs (found by review). In parquet-go 0.29.0 a byte array
column buffer's Reset keeps the scratch buffer it swaps in to write a page of
one value, so that page's bounds, when the value is empty, point into memory
a new writer does not have. A file now takes a new writer only when a UTF8
column holds a single non-null value and it is empty, or, as a precaution the
tests could not show is needed, when the file could span pages and holds an
empty string.

The parity tests hardly gave a reset writer a file: with a symbol in four
empty, almost every file took a new one. TestJSONParquetEncodeResetWriterWritesAsNew
encodes the same batches with a pooled writer and with a new one for every
file, with every encoding and codec, partitioned and not: one-row and few-row
files, empty, null and absent strings, and files spanning pages, ending in a
page of one empty string and clean of them. It requires the same bytes, and
that every file of a clean batch was written by a reset writer (with no GC to
empty the pool, and one P, whose private slot the writer is put in). Without
the check, 15 of its 36 configurations fail; without the one-value rule, 12.

A writer is reset to nil before it is pooled, so that it does not keep the
file it wrote reachable. The comments no longer say writers are kept between
batches, which the pool does not promise, and the docs give the cause of a
new writer.

Encode benchmark, 710 files: 117 -> 103 MB per batch, peak heap 101 -> 87 MiB.
26 files: unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The "escaped field name" case replaced "symbol" with "symbol", so no test sent
a name that matches a column only once unescaped: the path that reads a name
as raw JSON and unescapes it to look it up (found by review). It now sends
"symbol", and two more cases send escaped names with escaped values, and
an escaped name no column has before one that matches. All three fail when
names are looked up without being unescaped.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ssage

The encode benchmark split its batch into many files and few. It now also
writes it as one file, and as a file a message (json_parquet_encode only: the
legacy pipeline takes minutes), the shapes where the output buffers and the
arenas of small partitions cost most.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A partition's arena began with chunks of 8 rows and 1 KiB, held until its
file is written, so a batch split into partitions of a row or a few held
several times the memory its rows take (found by review). Chunks now begin at
one row and 64 bytes and double as before.

Encode benchmark, per batch: a file a message (35966 files) 1484 -> 1368 MB,
peak heap 277 -> 251 MiB; 710 files 101 -> 99 MB; 26 files and one file
unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A row's partition is found before its values are coerced, so a rejected row
of a new partition made its group, key string and arena chunks, and dropped
them (found by review). The group of a new partition is now a spare, given
its key and first message only once a row of it is kept, and a rejected row
gives back the chunks it began as well as its place in older ones.

BenchmarkJSONParquetEncodeRejected, 10000 rows lacking a column, each of a
partition of its own: 10.3 -> 6.9 MB and 240k -> 200k allocations per batch,
81-96 -> 65-73 ms. What is left is mostly each row's error, formatted and
logged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…them out

With the writer's buffering gone, each file was written in small pieces into
a buffer of its own, which grew by doubling from nothing and was handed to
the file's message with what doubling left unused (found by review). The
files of a batch are now written into one buffer, grown once to the largest
of them, and each is copied out at its size. The last keeps the buffer when
at most a quarter of it is unused, so a batch of one file is not held twice.

Encode benchmark, per batch: a file a message 1368 -> 1176 MB, peak heap
251 -> 215 MiB; 710 files 99 -> 94 MB, peak 85 -> 82 MiB; 26 files
73 -> 70 MB; one file unchanged. Time unchanged within noise (A/B run
back to back on 26 files: 356-369 against 361 ms).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
An interpolated column's string, which a cached one is for the whole batch,
was copied into the arena for every row, as was the string a bool or nested
value is coerced to (found by review). Both are strings no one writes to,
and the writer copies values into its own buffers, so a value now points at
the string's bytes. The arena copies only what it must: a field's content,
which the next message is parsed over. copy and copyString, one-line
wrappers of a generic, become one method.

Encode benchmark: 710 files 93.6 -> 92.2 MB per batch; one file, A/B back to
back, 95.3 -> 94.9 MB; time unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Every member name was scanned for an escape and for invalid UTF-8 before it
was looked up, though the decoder had just read and checked it (found by
review). Columns are now also kept by their names as JSON quotes them, so a
name is looked up as it is, and only one that is not so quoted and has an
escape or invalid UTF-8 is unquoted to be looked up again. A name quoted
without either cannot name a column otherwise: it would have to be that
name's own bytes, which is how JSON quotes it.

TestJSONParquetEncodeColumnLookup compares the lookup against unquoting every
name, for names with quotes, backslashes, <&>, U+2028, DEL, tabs and non-ASCII,
each as JSON quotes it, raw, fully \u-escaped and with invalid UTF-8; it and
the escaped-name parity cases fail when the unquoting is dropped.

Parsing a message, back to back: 2120-2570 -> 1800-1970 ns (-14%).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- A field's kind is jsontext.Value.Kind, rather than the first byte mapped
  by hand, and an absent field is kind 0, which drops jpeRow.present.
- A number's content is converted to a string for strconv as it is: the
  string does not escape, so one of up to 32 bytes is not allocated, and
  jpeUnsafeString goes.
- An arena's next chunk is sized from the last one's capacity, which drops
  its nextRows and nextBytes.
- unquoteString's comment names it.
- The tests make every encoder through newTestJSONEncoder, which takes
  resources and a testing.TB.

No change in what is written or allocated (0.91M, 0.76M and 0.72M
allocations per batch as before).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@almostintuitive almostintuitive changed the title parquet: json_parquet_encode allocates a third as often, and reuses its writers parquet: json_parquet_encode allocates a third as often, reuses writers, caches by metadata Sep 29, 2026
…stopping GC

TestJSONParquetEncodeResetWriterWritesAsNew stopped the GC, so that the
sync.Pool would keep its writers between batches, and pinned GOMAXPROCS to
one P. With the GC stopped, the garbage of 36 configurations of writers
(brotli's among them) grew to 13 GB, and CI's runner killed the job.

The encoder's pool is now an interface a sync.Pool satisfies, and the test
gives the encoder one that keeps what it is given. The test peaks at 128 MB,
the package at 351 MB. Without the check it guards, 12 of its 36
configurations still fail.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
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