parquet: json_parquet_encode allocates a third as often, reuses writers, caches by metadata - #5
Open
almostintuitive wants to merge 20 commits into
Open
almostintuitive wants to merge 20 commits into
almostintuitive wants to merge 20 commits into
Conversation
…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>
…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>
8 of 9 tasks
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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.strconvcall the Bloblang coercion makes; bools and nested values still go through the original helpers.columns[].cache_by(opt-in). An interpolated column is evaluated once per batch for each distinct combination of the listed metadata fields.json_parquet_encode_filescounter makes the file count visible in production.Review fixes
[]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: truecould 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 nowcache_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 publicXUnwrapperis gone; the expression is parsed through the manager's Bloblang environment, asgroup_by_valuedoes.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.count()test failed under-count=2, because Bloblang's counters are process-global; it now uses a new counter each run.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.
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_byis 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.Resetupstream would remove the need for new writers entirely. That isn't filed yet.Test plan
go test -count=2 ./internal/impl/parquet/, andgo test ./internal/bloblang/... ./public/service/ ./internal/impl/kafka/bento_docs_gen🤖 Generated with Claude Code