Skip to content

protobuf_parquet_encode: protobuf to Parquet in one processor - #1

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

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

Conversation

@almostintuitive

Copy link
Copy Markdown

Three commits, kept separate because only the first is of general interest.

protobuf_parquet_encode

A batch processor that turns protobuf-encoded messages straight into a Parquet file: protowire decode → Arrow builders → pqarrow, one row group per file. Configured with the message type, import_paths, and an explicit columns list (each a field of the message or a constant), so the Parquet schema is pinned by config rather than inferred.

partition splits one batch into several files by a timestamp field (field, unit, Go layout), emitting one message per partition key with the key exposed as metadata — the intended use is an S3 output path like .../${! meta("partition") }/..., with the encode placed inside the output's batching policy so the input's offsets are only committed once the upload is acknowledged.

skip_invalid_messages (default true) drops and counts undecodable messages via a protobuf_parquet_encode_invalid counter rather than failing the batch.

Benchmarked at ~3.2M msgs/s per core on the encode itself; end to end (Kafka → Parquet → S3) the pipeline runs ~1.04M msgs/s on 2 cores, against 148k for the same pipeline built from parquet_encode on stock Bento.

kafka_franz: add_record_metadata

Kafka record metadata is copied onto every message unconditionally. For a high-volume pipeline that never reads it, that allocation is pure cost, so this makes it optional (default true — existing behaviour unchanged).

Docker cross-compilation

resources/docker/Dockerfile builds with FROM --platform=$BUILDPLATFORM and honours TARGETOS/TARGETARCH, so multi-arch images build without emulation. Plus a workflow that publishes ghcr.io/aperiodic-io/bento for amd64 and arm64.

🤖 Generated with Claude Code

https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs

almostintuitive and others added 5 commits September 16, 2026 01:00
Decodes serialized protobuf messages straight from the wire format into
typed Arrow column builders and writes one single-row-group Parquet file
per batch (or per timestamp-derived partition key), instead of chaining
protobuf to_json, a Bloblang mapping and parquet_encode. On Kafka
archival workloads this removes the JSON round trip, per-message maps
and most GC pressure: ~7x throughput and ~9x msgs/s per core end to end,
~3M msgs/s per core for the processor itself.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs
Setting kafka_* metadata and headers costs several allocations per
record. Pipelines that never read them (e.g. archiving raw payloads) can
now opt out; the default keeps the existing behaviour.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs
Build on the native platform and cross-compile for the target so arm64
images don't compile Go under QEMU, and publish multi-arch images to
ghcr.io/aperiodic-io/bento from this fork.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs
The existing tests exercise the types our own pipelines use -- int64, double,
string, bool, packed repeated double -- which leaves most of the decoder's
branches unproven for anyone with a different schema. Each scalar kind has its
own decode path and its own Arrow builder, and several reinterpret bits
(zigzag, sign extension, float bit patterns) in ways the happy path never
touches.

Adds, against a message carrying every supported kind at the min/max of its
width:

* every scalar kind round-tripped, including the enum and a field number past
  15 (a two-byte tag);
* unpacked repeated fields -- legal protobuf that proto2 encoders and
  hand-rolled writers emit, and which a packed-only decoder reads as empty;
* unknown fields of all four wire types skipped ahead of a mapped one, since
  consuming the wrong length corrupts every field that follows;
* last-value-wins for a singular field repeated on the wire;
* all four partition units dividing one instant to the same day;
* the compression options asserted against the Parquet metadata, which is the
  only place they are observable -- a file written with the wrong codec reads
  back identically;
* the repeated-string and map guards, the first of which is load-bearing: a
  repeated string arrives as one length-delimited chunk per element, which the
  packed-scalar path would misread as a run of varints.

Verified by mutation: inverting the sfixed32 reinterpret and the ms divisor
each fail exactly one of the new assertions.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs
…radicting its row

Review findings on protobuf_parquet_encode, each reproduced by a test that fails
without the fix:

* An out-of-enum `partition.unit` reached the constructor intact -- the enum on
  the field is a lint rule, not a parse-time constraint, so anything that skips
  linting got divisor 0 and a panic (integer divide by zero) on the first
  message. Now a config error.
* An unknown `compression` resolved to the zero value, which is Uncompressed, so
  a typo silently wrote uncompressed Parquet forever. Now a config error.
* The partition key and the column value are derived from the same bits by two
  paths that disagreed on sfixed32: a row whose own timestamp column read
  -86400 was filed under 2106-02-06 instead of 1969-12-31. Both now go through
  one ppeSignedInt helper.
* The slot table is indexed by field number, so the highest number used sized it
  -- 8MB for a single column on field 1,000,000, gigabytes near protobuf's limit
  of 536,870,911. Past 4096 it now uses a map.
* Every partition group reserved capacity for the whole batch, multiplying peak
  memory by the number of keys a batch straddles. It now hints only the rows
  that can still arrive.
* The sint32 appender boxed an int32 into an `any` and asserted it back on every
  row -- the one allocation left on the hot path. Inlined; ppeScalar is gone.

Not changed: a length-delimited payload on a packable repeated field decoding as
packed varints. That is the wire format, not a bug -- proto.Unmarshal returns
exactly the same values for the same bytes, and a test now pins that parity so
the decoder is not "fixed" into disagreeing with the protobuf runtime.

Benchmark is unchanged (3.5M msgs/s per core).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs
@almostintuitive
almostintuitive added this pull request to stack #4 September 29, 2026 11:39
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