Skip to content

parquet: json_parquet_encode; kafka_franz: copy_record_values - #3

Open
almostintuitive wants to merge 5 commits into
perf/protobuf-parquet-encodefrom
perf/json-parquet-encode
Open

almostintuitive wants to merge 5 commits into
perf/protobuf-parquet-encodefrom
perf/json-parquet-encode

Conversation

@almostintuitive

@almostintuitive almostintuitive commented Sep 29, 2026 •

Copy link
Copy Markdown

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:

    • a per-column Bloblang mapping (.not_null().string(), .int32(), .int64(), .number() / .or(NaN).number());
    • a catch that drops rejected rows;
    • group_by_value on a partition path;
    • parquet_encode.

    How it guarantees that:

    • the schema and writer are parquet_encode's, built by the same code (GenerateStructType, parquet.SchemaOf, GenericWriter);
    • each column is coerced by the same internal/value helpers the Bloblang methods call (IToString, IToInt32, IToInt, IToNumber);
    • coercion runs on the Go values Bento's own JSON decoding produces (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 optional partition.path ({column} and {column|go-time-layout} placeholders) splits the batch into one file per path.

  • kafka_franz.copy_record_values (advanced, default false). 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.go and processor_json_encode_parity_cases_test.go run the replaced pipeline and json_parquet_encode as 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.

  • Every column type × optionality against 68 values: every number literal form and range edge, numeric strings (hex, octal, underscores, 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.
  • Malformed messages: empty, whitespace, truncated, arrays and scalars, two documents, trailing garbage, BOM, duplicate keys, deep nesting, 1 MiB strings.
  • Metadata: missing topic, short topic, empty segment, unicode.
  • Partitions: around midnight, before the epoch, far future, fractional, string and overflowing times; path fields of every JSON kind.
  • Random batches: 300 seeded batches (60,313 messages, 3,137 dropped), plus a 10 MiB batch (40,949 messages across 747 files).
  • Fuzz: FuzzJSONParquetParity ran 15 minutes locally, 8.1M executions and 1,151 interesting inputs, with no divergence.
  • Mutation check: 13 hand-made mutants of the processor. 12 were caught; the 13th showed a redundant check, which was removed.

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):

legacy pipeline json_parquet_encode
encode time 2.30 s 0.66 s (3.5x faster)
allocations 12.4M 2.6M
bytes allocated 603 MB 201 MB
held while batching (per byte of JSON) 5.15 B 2.93 B
peak heap during encode 230 MiB 101 MiB

Rows are collected per partition, and each partition's file is encoded and released before the next one starts, the same way parquet_encode handles each group from group_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 through Reset doesn't write the same bytes; the parity tests catch that.

Test plan

  • go test ./internal/impl/parquet/ ./internal/impl/kafka/
  • go vet on both packages
  • fuzz, 15 minutes
  • make docs: component docs regenerated. The lint step fails locally on the DuckDB/sql_raw docs, which need a CGO build; unrelated to this change.
  • golangci-lint v2.13.2 (built with Go 1.27.1) on both packages: 0 issues
  • fuzz, 90 seconds after the one-file-at-a-time change: 2.1M executions, no divergence
  • CI: the new Aperiodic Test workflow (.github/workflows/aperiodic_test.yml) runs vet, tests and golangci-lint on ubicloud-standard-2. It covers only the Go packages the PR touches (the nearest directory with Go files to each changed file) and checks that go.mod/go.sum are tidy when they change. Upstream's Test workflow is still disabled on this fork.

🤖 Generated with Claude Code

almostintuitive and others added 3 commits September 29, 2026 13:15
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
almostintuitive added this pull request to stack #4 September 29, 2026 11:39
almostintuitive and others added 2 commits September 29, 2026 11:53
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>
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