protobuf_parquet_encode: protobuf to Parquet in one processor - #1
Open
almostintuitive wants to merge 5 commits into
Open
almostintuitive wants to merge 5 commits into
almostintuitive wants to merge 5 commits into
Conversation
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
added this pull request to stack #4
September 29, 2026 11:39
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.
Three commits, kept separate because only the first is of general interest.
protobuf_parquet_encodeA 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 explicitcolumnslist (each afieldof the message or aconstant), so the Parquet schema is pinned by config rather than inferred.partitionsplits one batch into several files by a timestamp field (field,unit, Golayout), 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 aprotobuf_parquet_encode_invalidcounter 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_encodeon stock Bento.kafka_franz: add_record_metadataKafka 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/Dockerfilebuilds withFROM --platform=$BUILDPLATFORMand honoursTARGETOS/TARGETARCH, so multi-arch images build without emulation. Plus a workflow that publishesghcr.io/aperiodic-io/bentofor amd64 and arm64.🤖 Generated with Claude Code
https://claude.ai/code/session_017oaFUVie2SYfkDdnrRk9bs