parquet: json_parquet_encode; kafka_franz: copy_record_values - #3
Open
almostintuitive wants to merge 5 commits into
Open
almostintuitive wants to merge 5 commits into
almostintuitive wants to merge 5 commits into
Conversation
franz-go hands out records whose key, value and header values are slices of the fetch response they arrived in, which holds every record fetched from that broker in the same request. A message held for long (in a large output batch, say) therefore keeps the whole response in memory, and so does the record kept for checkpointing through its headers. With copy_record_values the message owns a copy of its value, and the record drops its key, value and headers once metadata is extracted. Off by default. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Encodes a batch of flat JSON objects straight into Parquet, row by row, and writes what a per-column Bloblang mapping, a catch that drops rejected rows, group_by_value and parquet_encode would: the schema and writer are parquet_encode's (built by the same code), and each column is coerced by the very value helpers Bloblang's .string(), .int32(), .int64() and .number() use, on the values Bento's own JSON decoding gives. A batch is therefore held only as its messages' bytes, never as structured rows. Rows a column rejects are dropped, logged and counted (json_parquet_encode_dropped). An optional partition path splits the batch into one file per path, formatted from row fields and a time layout. The parity tests run the replaced pipeline and json_parquet_encode as real streams over the same batches (every column type and optionality against ~70 values, malformed messages, metadata, partition edge cases, 300 random batches, a 10 MiB batch) and require byte-identical files, partitions and metadata; a fuzz target does the same. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… modernize) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
almostintuitive
added this pull request to stack #4
September 29, 2026 11:39
A batch spanning many partitions held one writer, with its column buffers, per partition until every row was read, so its peak heap during encoding grew with the partition count and outgrew the pipeline it replaces (~300 MiB against ~230 MiB on a 10 MiB batch across 747 partitions). Rows are now collected per partition and each file is encoded and released before the next, as parquet_encode does for the groups group_by_value hands it: the same batch peaks at ~100 MiB and allocates 200 MB instead of 356 MB. Each file takes a new writer: a writer reset with Reset does not write the same bytes as a new one, which the parity tests catch. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The module proxy intermittently fails a download with a stream error, which failed the first run in go vet. Download what the touched packages need first, retrying. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
5 tasks done
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
Two changes that let a JSON→Parquet archiver hold its batches cheaply:
json_parquet_encode(new batch processor,internal/impl/parquet/processor_json_encode.go). It encodes flat JSON objects straight into Parquet, row by row, and writes the same bytes as the pipeline it replaces:.not_null().string(),.int32(),.int64(),.number()/.or(NaN).number());catchthat drops rejected rows;group_by_valueon a partition path;parquet_encode.How it guarantees that:
parquet_encode's, built by the same code (GenerateStructType,parquet.SchemaOf,GenericWriter);internal/valuehelpers the Bloblang methods call (IToString,IToInt32,IToInt,IToNumber);json.Number, last duplicate key wins, U+FFFD for invalid UTF-8, exactly one document).Rows it rejects are dropped, logged, and counted in
json_parquet_encode_dropped. An optionalpartition.path({column}and{column|go-time-layout}placeholders) splits the batch into one file per path.kafka_franz.copy_record_values(advanced, defaultfalse). Copies each record's value into its message, and drops the record's key, value and headers once metadata is extracted. Without it, a message held for a long time (for example in a large output batch) keeps its whole fetch response in memory, along with every other record fetched from that broker in the same request.Parity tests
processor_json_encode_parity_test.goandprocessor_json_encode_parity_cases_test.gorun the replaced pipeline andjson_parquet_encodeas real in-process streams over the same batches. They require byte-identical Parquet files, the same partitions, the same metadata, and the same rows dropped.NaN,inf), unicode, invalid UTF-8, escapes, bools, nulls, nested values, and missing fields. Each column is asserted to have both kept and dropped values, so a case can't pass by both sides writing nothing.FuzzJSONParquetParityran 15 minutes locally, 8.1M executions and 1,151 interesting inputs, with no divergence.The one intentional divergence: an optional UTF8 column reads null or absent as NULL. The mapping it's modelled on would have rejected that row. This is covered by
TestJSONParquetEncodeOptionalStringIsNull.Benchmarks
AMD EPYC (16 vCPU), 10 MiB batch of synthetic rows across 710 partitions (
processor_json_encode_bench_test.go):Rows are collected per partition, and each partition's file is encoded and released before the next one starts, the same way
parquet_encodehandles each group fromgroup_by_value. An earlier revision kept one writer per partition open until the whole batch was read. On this batch that meant 710 writers with their column buffers, and a 300 MiB peak, more than the legacy path. Each file gets a new writer, because a writer reused throughResetdoesn't write the same bytes; the parity tests catch that.Test plan
go test ./internal/impl/parquet/ ./internal/impl/kafka/go veton both packagesmake docs: component docs regenerated. The lint step fails locally on the DuckDB/sql_raw docs, which need a CGO build; unrelated to this change.Aperiodic Testworkflow (.github/workflows/aperiodic_test.yml) runs vet, tests and golangci-lint onubicloud-standard-2. It covers only the Go packages the PR touches (the nearest directory with Go files to each changed file) and checks thatgo.mod/go.sumare tidy when they change. Upstream'sTestworkflow is still disabled on this fork.🤖 Generated with Claude Code