diff --git a/.github/workflows/aperiodic_image.yml b/.github/workflows/aperiodic_image.yml new file mode 100644 index 0000000000..6b54107d88 --- /dev/null +++ b/.github/workflows/aperiodic_image.yml @@ -0,0 +1,41 @@ +name: Aperiodic Image + +on: + push: + branches: [main, 'perf/**'] + workflow_dispatch: {} + +permissions: + contents: read + packages: write + +jobs: + image: + if: github.repository == 'aperiodic-io/bento' + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v7 + + - uses: docker/setup-buildx-action@v4 + + - uses: docker/login-action@v4 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + + - name: Compute tags + id: tags + run: | + branch="$(echo "${GITHUB_REF_NAME}" | tr '/' '-' | tr -cd 'A-Za-z0-9_.-')" + echo "tags=ghcr.io/${{ github.repository_owner }}/bento:${GITHUB_SHA::10},ghcr.io/${{ github.repository_owner }}/bento:${branch}" >> "$GITHUB_OUTPUT" + + - uses: docker/build-push-action@v7 + with: + context: ./ + file: ./resources/docker/Dockerfile + platforms: linux/amd64,linux/arm64 + push: true + tags: ${{ steps.tags.outputs.tags }} + cache-from: type=gha + cache-to: type=gha,mode=max diff --git a/go.mod b/go.mod index 9d6b305389..4427418da0 100644 --- a/go.mod +++ b/go.mod @@ -32,6 +32,7 @@ require ( github.com/PaesslerAG/gval v1.2.3 github.com/PaesslerAG/jsonpath v0.1.1 github.com/andybalholm/brotli v1.2.3 + github.com/apache/arrow-go/v18 v18.8.0 github.com/apache/pulsar-client-go v0.21.0 github.com/aws/aws-lambda-go v1.46.0 github.com/aws/aws-msk-iam-sasl-signer-go v1.0.4 @@ -224,8 +225,8 @@ require ( github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 // indirect github.com/RoaringBitmap/roaring/v2 v2.8.0 // indirect github.com/alexbrainman/sspi v0.0.0-20250919150558-7d374ff0d59e // indirect - github.com/apache/arrow-go/v18 v18.8.0 // indirect github.com/apache/arrow/go/v15 v15.0.2 // indirect + github.com/apache/thrift v0.24.0 // indirect github.com/ardielle/ardielle-go v1.5.2 // indirect github.com/armon/go-metrics v0.3.4 // indirect github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 // indirect diff --git a/internal/impl/kafka/input_kafka_franz.go b/internal/impl/kafka/input_kafka_franz.go index 5ec4efdc7a..2da0ceca9f 100644 --- a/internal/impl/kafka/input_kafka_franz.go +++ b/internal/impl/kafka/input_kafka_franz.go @@ -166,6 +166,7 @@ With this option, you can return topic order and per-topic partition ordering. T Field(service.NewTLSToggledField("tls")). Field(saslField()). Field(service.NewBoolField("multi_header").Description("Decode headers into lists to allow handling of multiple values with the same key").Default(false).Advanced()). + Field(service.NewBoolField("add_record_metadata").Description("Add the `kafka_*` metadata fields and record headers to each message. Disabling this avoids several allocations per record when no downstream component reads them.").Default(true).Advanced()). Field(service.NewBatchPolicyField("batching"). Description("Allows you to configure a [batching policy](/docs/configuration/batching) that applies to individual topic partitions in order to batch messages together before flushing them for processing. Batching can be beneficial for performance as well as useful for windowed processing, and doing so this way preserves the ordering of topic partitions."). Advanced()). @@ -227,6 +228,7 @@ type franzKafkaReader struct { commitPeriod time.Duration regexPattern bool multiHeader bool + addMetadata bool batchPolicy service.BatchPolicy reconnectOnUnknownTopic bool @@ -500,6 +502,9 @@ func newFranzKafkaReaderFromConfig(conf *service.ParsedConfig, res *service.Reso if tlsEnabled { f.tlsConf = tlsConf } + if f.addMetadata, err = conf.FieldBool("add_record_metadata"); err != nil { + return nil, err + } if f.multiHeader, err = conf.FieldBool("multi_header"); err != nil { return nil, err } @@ -517,6 +522,11 @@ type msgWithRecord struct { func (f *franzKafkaReader) recordToMessage(record *kgo.Record) *msgWithRecord { msg := service.NewMessage(record.Value) + if !f.addMetadata { + record.Key = nil + record.Value = nil + return &msgWithRecord{msg: msg, r: record} + } if record.Key != nil { msg.MetaSetMut("kafka_key", string(record.Key)) } diff --git a/internal/impl/kafka/input_kafka_franz_metadata_test.go b/internal/impl/kafka/input_kafka_franz_metadata_test.go new file mode 100644 index 0000000000..d7e7eb0841 --- /dev/null +++ b/internal/impl/kafka/input_kafka_franz_metadata_test.go @@ -0,0 +1,41 @@ +package kafka + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/twmb/franz-go/pkg/kgo" +) + +func TestFranzRecordToMessageMetadata(t *testing.T) { + newRecord := func() *kgo.Record { + return &kgo.Record{ + Key: []byte("k"), Value: []byte("payload"), Topic: "raw.quotes", Partition: 1, Offset: 42, + Timestamp: time.Unix(1_789_502_374, 0), Headers: []kgo.RecordHeader{{Key: "redpanda-dedup-key", Value: []byte("abc")}}, + } + } + + withMeta := (&franzKafkaReader{addMetadata: true}).recordToMessage(newRecord()) + b, err := withMeta.msg.AsBytes() + require.NoError(t, err) + assert.Equal(t, "payload", string(b)) + for key, want := range map[string]any{"kafka_key": "k", "kafka_topic": "raw.quotes", "kafka_partition": 1, "kafka_offset": 42, "redpanda-dedup-key": "abc"} { + got, ok := withMeta.msg.MetaGetMut(key) + require.True(t, ok, key) + assert.Equal(t, want, got, key) + } + + // Opting out must keep the payload and the record (needed for offset + // checkpointing) while adding no metadata at all. + bare := (&franzKafkaReader{addMetadata: false}).recordToMessage(newRecord()) + b, err = bare.msg.AsBytes() + require.NoError(t, err) + assert.Equal(t, "payload", string(b)) + assert.Equal(t, int64(42), bare.r.Offset) + assert.Equal(t, int32(1), bare.r.Partition) + count := 0 + require.NoError(t, bare.msg.MetaWalkMut(func(string, any) error { count++; return nil })) + assert.Zero(t, count) +} diff --git a/internal/impl/protobuf/processor_protobuf_parquet.go b/internal/impl/protobuf/processor_protobuf_parquet.go new file mode 100644 index 0000000000..1220affde7 --- /dev/null +++ b/internal/impl/protobuf/processor_protobuf_parquet.go @@ -0,0 +1,699 @@ +package protobuf + +import ( + "bytes" + "context" + "errors" + "fmt" + "math" + "time" + + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/arrow-go/v18/parquet" + "github.com/apache/arrow-go/v18/parquet/compress" + "github.com/apache/arrow-go/v18/parquet/pqarrow" + "google.golang.org/protobuf/encoding/protowire" + "google.golang.org/protobuf/reflect/protoreflect" + + "github.com/warpstreamlabs/bento/public/service" +) + +const ( + ppeFieldMessage = "message" + ppeFieldImportPaths = "import_paths" + ppeFieldColumns = "columns" + ppeFieldColumnName = "name" + ppeFieldColumnField = "field" + ppeFieldColumnConstant = "constant" + ppeFieldPartition = "partition" + ppeFieldPartitionField = "field" + ppeFieldPartitionUnit = "unit" + ppeFieldPartitionLayout = "layout" + ppeFieldPartitionMetaKey = "metadata_key" + ppeFieldCompression = "compression" + ppeFieldCompressionLevel = "compression_level" + ppeFieldDictionary = "dictionary" + ppeFieldSkipInvalid = "skip_invalid_messages" +) + +func protobufParquetEncodeSpec() *service.ConfigSpec { + return service.NewConfigSpec(). + Beta(). + Categories("Parsing"). + Summary("Decodes a batch of protobuf messages straight into columnar Parquet files, without an intermediate JSON or structured representation."). + Description(` +Each message in the batch must be a single serialized protobuf message of the configured type. The wire format is decoded directly into typed Apache Arrow column builders and written as a Parquet file with one row group, which is an order of magnitude cheaper than chaining the `+"`protobuf`"+` (to_json), `+"`mapping`"+` and `+"`parquet_encode`"+` processors. + +Columns are emitted in the configured order. A column either copies a protobuf field (optionally renamed) or holds a constant string. Proto3 fields that are absent from a message take their zero value. Repeated scalar fields become Parquet LIST columns and string columns are dictionary encoded when `+"`dictionary`"+` is enabled. + +When `+"`partition`"+` is set, the batch is split by formatting a timestamp field with a Go time layout, producing one Parquet file per distinct key, and the key is written to the configured metadata field of each output message. Every output message otherwise carries the metadata of the first input message of its partition. + +Messages that cannot be decoded are dropped and logged when `+"`skip_invalid_messages`"+` is true, and fail the batch otherwise.`). + Field(service.NewStringField(ppeFieldMessage).Description("The fully qualified name of the protobuf message.")). + Field(service.NewStringListField(ppeFieldImportPaths).Description("Directories to load `.proto` files from.").Default([]any{})). + Field(service.NewObjectListField(ppeFieldColumns, + service.NewStringField(ppeFieldColumnName).Description("The Parquet column name."), + service.NewStringField(ppeFieldColumnField).Description("The protobuf field copied into the column. Defaults to the column name.").Optional(), + service.NewStringField(ppeFieldColumnConstant).Description("A constant string value for every row, instead of a protobuf field.").Optional(), + ).Description("The output columns, in order.")). + Field(service.NewObjectField(ppeFieldPartition, + service.NewStringField(ppeFieldPartitionField).Description("An integer protobuf field holding a Unix timestamp."), + service.NewStringEnumField(ppeFieldPartitionUnit, "s", "ms", "us", "ns").Description("The unit of the timestamp field.").Default("us"), + service.NewStringField(ppeFieldPartitionLayout).Description("A Go time layout, evaluated in UTC, that forms the partition key.").Example("year=2006/month=01/day=02/hour=15"), + service.NewStringField(ppeFieldPartitionMetaKey).Description("The metadata field that receives the partition key.").Default("partition"), + ).Description("Split the batch into one Parquet file per formatted timestamp key.").Optional()). + Field(service.NewStringEnumField(ppeFieldCompression, "zstd", "snappy", "gzip", "lz4raw", "uncompressed").Default("zstd")). + Field(service.NewIntField(ppeFieldCompressionLevel).Description("The compression level, for codecs that support one.").Default(3).Advanced()). + Field(service.NewBoolField(ppeFieldDictionary).Description("Dictionary encode string columns.").Default(true).Advanced()). + Field(service.NewBoolField(ppeFieldSkipInvalid).Description("Drop (and log) messages that cannot be decoded instead of failing the batch.").Default(true)). + Example("Archiving protobuf records to S3 as hourly Parquet", + "Batches records at the output, splits each batch by hour of a microsecond timestamp and uploads one Parquet file per hour.", + ` +output: + aws_s3: + bucket: raw-archive + path: 'quotes/${! meta("partition") }/${! uuid_v4() }.parquet' + batching: + count: 50000 + period: 60s + processors: + - protobuf_parquet_encode: + message: live.producers.v1.QuoteRecord + import_paths: [ /etc/bento/proto ] + columns: + - { name: exchange, constant: binance-futures } + - { name: symbol } + - { name: time } + - { name: local_timestamp } + - { name: bid_price } + partition: + field: local_timestamp + unit: us + layout: year=2006/month=01/day=02/hour=15 +`) +} + +func init() { + err := service.RegisterBatchProcessor("protobuf_parquet_encode", protobufParquetEncodeSpec(), + func(conf *service.ParsedConfig, mgr *service.Resources) (service.BatchProcessor, error) { + return newProtobufParquetEncoder(conf, mgr) + }) + if err != nil { + panic(err) + } +} + +// ------------------------------------------------------------------------------ + +// ppeColumn is one output column. slot indexes the decoded value of the row, or +// is -1 for constant columns. +type ppeColumn struct { + name string + constant string + slot int +} + +// ppeSlot is one protobuf field decoded per row. +type ppeSlot struct { + kind protoreflect.Kind + repeated bool + dataType arrow.DataType +} + +type protobufParquetEncoder struct { + log *service.Logger + mInvalid *service.MetricCounter + columns []ppeColumn + slots []ppeSlot + slotByField []int // indexed by protobuf field number, -1 when not decoded + slotOfField map[protoreflect.FieldNumber]int // used instead when field numbers are sparse + schema *arrow.Schema + props *parquet.WriterProperties + skipInvalid bool + + partitionSlot int // -1 when partitioning is disabled + partitionDivisor int64 + partitionLayout string + partitionMetaKey string +} + +func newProtobufParquetEncoder(conf *service.ParsedConfig, mgr *service.Resources) (*protobufParquetEncoder, error) { + msgName, err := conf.FieldString(ppeFieldMessage) + if err != nil { + return nil, err + } + importPaths, err := conf.FieldStringList(ppeFieldImportPaths) + if err != nil { + return nil, err + } + files, _, err := loadDescriptors(mgr.FS(), importPaths) + if err != nil { + return nil, fmt.Errorf("failed to load protobuf definitions: %w", err) + } + d, err := files.FindDescriptorByName(protoreflect.FullName(msgName)) + if err != nil { + return nil, fmt.Errorf("unable to find message '%v' definition: %w", msgName, err) + } + md, ok := d.(protoreflect.MessageDescriptor) + if !ok { + return nil, fmt.Errorf("'%v' is not a message", msgName) + } + + e := &protobufParquetEncoder{ + log: mgr.Logger(), + mInvalid: mgr.Metrics().NewCounter("protobuf_parquet_encode_invalid"), + partitionSlot: -1, + } + slotOfField := map[protoreflect.FieldNumber]int{} + addSlot := func(fieldName string) (int, error) { + fd := md.Fields().ByName(protoreflect.Name(fieldName)) + if fd == nil { + return 0, fmt.Errorf("message '%v' has no field '%v'", msgName, fieldName) + } + if s, exists := slotOfField[fd.Number()]; exists { + return s, nil + } + dt, err := ppeArrowType(fd) + if err != nil { + return 0, err + } + e.slots = append(e.slots, ppeSlot{kind: fd.Kind(), repeated: fd.IsList(), dataType: dt}) + slotOfField[fd.Number()] = len(e.slots) - 1 + return len(e.slots) - 1, nil + } + + colConfs, err := conf.FieldObjectList(ppeFieldColumns) + if err != nil { + return nil, err + } + if len(colConfs) == 0 { + return nil, errors.New("at least one column is required") + } + arrowFields := make([]arrow.Field, 0, len(colConfs)) + for _, cc := range colConfs { + name, err := cc.FieldString(ppeFieldColumnName) + if err != nil { + return nil, err + } + col := ppeColumn{name: name, slot: -1} + if cc.Contains(ppeFieldColumnConstant) { + if cc.Contains(ppeFieldColumnField) { + return nil, fmt.Errorf("column '%v' sets both field and constant", name) + } + if col.constant, err = cc.FieldString(ppeFieldColumnConstant); err != nil { + return nil, err + } + arrowFields = append(arrowFields, arrow.Field{Name: name, Type: arrow.BinaryTypes.String}) + } else { + fieldName := name + if cc.Contains(ppeFieldColumnField) { + if fieldName, err = cc.FieldString(ppeFieldColumnField); err != nil { + return nil, err + } + } + if col.slot, err = addSlot(fieldName); err != nil { + return nil, fmt.Errorf("column '%v': %w", name, err) + } + arrowFields = append(arrowFields, arrow.Field{Name: name, Type: e.slots[col.slot].dataType}) + } + e.columns = append(e.columns, col) + } + e.schema = arrow.NewSchema(arrowFields, nil) + + if conf.Contains(ppeFieldPartition) { + pc := conf.Namespace(ppeFieldPartition) + fieldName, err := pc.FieldString(ppeFieldPartitionField) + if err != nil { + return nil, err + } + if e.partitionSlot, err = addSlot(fieldName); err != nil { + return nil, fmt.Errorf("partition: %w", err) + } + if s := e.slots[e.partitionSlot]; s.repeated || !ppeIsInteger(s.kind) { + return nil, fmt.Errorf("partition field '%v' must be a singular integer field", fieldName) + } + unit, err := pc.FieldString(ppeFieldPartitionUnit) + if err != nil { + return nil, err + } + // The enum on the field is a lint rule, not a parse-time constraint, so an + // unknown unit arrives here intact whenever linting is skipped. Left + // unchecked it becomes divisor 0 and panics on the first message. + divisor, ok := map[string]int64{"s": 1, "ms": 1e3, "us": 1e6, "ns": 1e9}[unit] + if !ok { + return nil, fmt.Errorf("partition: unknown unit '%v'", unit) + } + e.partitionDivisor = divisor + if e.partitionLayout, err = pc.FieldString(ppeFieldPartitionLayout); err != nil { + return nil, err + } + if e.partitionMetaKey, err = pc.FieldString(ppeFieldPartitionMetaKey); err != nil { + return nil, err + } + } + + maxField := protoreflect.FieldNumber(0) + for f := range slotOfField { + maxField = max(maxField, f) + } + // Indexed by field number, so it is only worth building while the highest + // number stays close to the number of fields. Protobuf allows numbers up to + // 536,870,911, and a schema that puts one column up there would otherwise + // allocate gigabytes here. + if maxField <= ppeMaxDenseFieldNumber { + e.slotByField = make([]int, maxField+1) + for i := range e.slotByField { + e.slotByField[i] = -1 + } + for f, s := range slotOfField { + e.slotByField[f] = s + } + } else { + e.slotOfField = slotOfField + } + + codecName, err := conf.FieldString(ppeFieldCompression) + if err != nil { + return nil, err + } + level, err := conf.FieldInt(ppeFieldCompressionLevel) + if err != nil { + return nil, err + } + dictionary, err := conf.FieldBool(ppeFieldDictionary) + if err != nil { + return nil, err + } + if e.skipInvalid, err = conf.FieldBool(ppeFieldSkipInvalid); err != nil { + return nil, err + } + // Unchecked, an unknown codec resolves to the zero value -- which is + // Uncompressed -- and silently writes uncompressed files forever. + codec, ok := map[string]compress.Compression{ + "zstd": compress.Codecs.Zstd, "snappy": compress.Codecs.Snappy, "gzip": compress.Codecs.Gzip, + "lz4raw": compress.Codecs.Lz4Raw, "uncompressed": compress.Codecs.Uncompressed, + }[codecName] + if !ok { + return nil, fmt.Errorf("unknown compression '%v'", codecName) + } + propOpts := []parquet.WriterProperty{ + parquet.WithCompression(codec), + parquet.WithDictionaryDefault(false), + // One row group per file: a batch is already the unit of upload. + parquet.WithMaxRowGroupLength(math.MaxInt64), + parquet.WithStats(true), + } + if codecName == "zstd" || codecName == "gzip" { + propOpts = append(propOpts, parquet.WithCompressionLevel(level)) + } + if dictionary { + for _, f := range arrowFields { + if f.Type.ID() == arrow.STRING { + propOpts = append(propOpts, parquet.WithDictionaryFor(f.Name, true)) + } + } + } + e.props = parquet.NewWriterProperties(propOpts...) + return e, nil +} + +// Past this, a slot table indexed by field number costs more than the map lookup +// it saves: 4096 entries is 32KB, a schema numbered beyond it is not dense. +const ppeMaxDenseFieldNumber = protoreflect.FieldNumber(4096) + +func ppeIsInteger(k protoreflect.Kind) bool { + switch k { + case protoreflect.Int64Kind, protoreflect.Sint64Kind, protoreflect.Sfixed64Kind, + protoreflect.Int32Kind, protoreflect.Sint32Kind, protoreflect.Sfixed32Kind, + protoreflect.Uint64Kind, protoreflect.Fixed64Kind, protoreflect.Uint32Kind, protoreflect.Fixed32Kind: + return true + } + return false +} + +func ppeArrowType(fd protoreflect.FieldDescriptor) (arrow.DataType, error) { + if fd.IsMap() { + return nil, fmt.Errorf("map field '%v' is not supported", fd.Name()) + } + var t arrow.DataType + switch fd.Kind() { + case protoreflect.Int64Kind, protoreflect.Sint64Kind, protoreflect.Sfixed64Kind: + t = arrow.PrimitiveTypes.Int64 + case protoreflect.Int32Kind, protoreflect.Sint32Kind, protoreflect.Sfixed32Kind, protoreflect.EnumKind: + t = arrow.PrimitiveTypes.Int32 + case protoreflect.Uint64Kind, protoreflect.Fixed64Kind: + t = arrow.PrimitiveTypes.Uint64 + case protoreflect.Uint32Kind, protoreflect.Fixed32Kind: + t = arrow.PrimitiveTypes.Uint32 + case protoreflect.DoubleKind: + t = arrow.PrimitiveTypes.Float64 + case protoreflect.FloatKind: + t = arrow.PrimitiveTypes.Float32 + case protoreflect.BoolKind: + t = arrow.FixedWidthTypes.Boolean + case protoreflect.StringKind: + t = arrow.BinaryTypes.String + case protoreflect.BytesKind: + t = arrow.BinaryTypes.Binary + default: + return nil, fmt.Errorf("field '%v' of kind %v is not supported", fd.Name(), fd.Kind()) + } + if fd.IsList() { + if t.ID() == arrow.STRING || t.ID() == arrow.BINARY { + return nil, fmt.Errorf("repeated %v field '%v' is not supported", fd.Kind(), fd.Name()) + } + return arrow.ListOf(t), nil + } + return t, nil +} + +// ------------------------------------------------------------------------------ + +// ppeRow holds the decoded values of one message, reused across messages. +// Scalars are kept as raw wire bits and interpreted per kind when appended. +type ppeRow struct { + bits []uint64 + bytes [][]byte + lists [][]uint64 +} + +func (r *ppeRow) reset() { + clear(r.bits) + clear(r.bytes) + for i := range r.lists { + r.lists[i] = r.lists[i][:0] + } +} + +func (e *protobufParquetEncoder) decode(b []byte, r *ppeRow) error { + r.reset() + for len(b) > 0 { + num, typ, n := protowire.ConsumeTag(b) + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + slot := -1 + if e.slotOfField != nil { + if mapped, ok := e.slotOfField[num]; ok { + slot = mapped + } + } else if int(num) < len(e.slotByField) { + slot = e.slotByField[num] + } + if slot < 0 { + if n = protowire.ConsumeFieldValue(num, typ, b); n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + continue + } + s := &e.slots[slot] + switch typ { + case protowire.VarintType: + v, n := protowire.ConsumeVarint(b) + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + e.store(r, slot, s, v) + case protowire.Fixed64Type: + v, n := protowire.ConsumeFixed64(b) + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + e.store(r, slot, s, v) + case protowire.Fixed32Type: + v, n := protowire.ConsumeFixed32(b) + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + e.store(r, slot, s, uint64(v)) + case protowire.BytesType: + v, n := protowire.ConsumeBytes(b) + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + if !s.repeated { + r.bytes[slot] = v + continue + } + if err := e.storePacked(r, slot, s, v); err != nil { + return err + } + default: + if n = protowire.ConsumeFieldValue(num, typ, b); n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + } + } + return nil +} + +func (e *protobufParquetEncoder) store(r *ppeRow, slot int, s *ppeSlot, v uint64) { + if s.repeated { + r.lists[slot] = append(r.lists[slot], v) + } else { + r.bits[slot] = v + } +} + +func (e *protobufParquetEncoder) storePacked(r *ppeRow, slot int, s *ppeSlot, b []byte) error { + for len(b) > 0 { + var v uint64 + var n int + switch s.kind { + case protoreflect.DoubleKind, protoreflect.Fixed64Kind, protoreflect.Sfixed64Kind: + v, n = protowire.ConsumeFixed64(b) + case protoreflect.FloatKind, protoreflect.Fixed32Kind, protoreflect.Sfixed32Kind: + var v32 uint32 + v32, n = protowire.ConsumeFixed32(b) + v = uint64(v32) + default: + v, n = protowire.ConsumeVarint(b) + } + if n < 0 { + return protowire.ParseError(n) + } + b = b[n:] + r.lists[slot] = append(r.lists[slot], v) + } + return nil +} + +// ------------------------------------------------------------------------------ + +type ppeGroup struct { + first *service.Message + key string + rows int + builders []array.Builder + appenders []func(r *ppeRow) +} + +func (e *protobufParquetEncoder) newGroup(first *service.Message, key string, mem memory.Allocator, sizeHint int) *ppeGroup { + g := &ppeGroup{first: first, key: key} + for _, col := range e.columns { + if col.slot < 0 { + bld := array.NewStringBuilder(mem) + bld.Reserve(sizeHint) + constant := col.constant + g.builders = append(g.builders, bld) + g.appenders = append(g.appenders, func(*ppeRow) { bld.Append(constant) }) + continue + } + bld := array.NewBuilder(mem, e.slots[col.slot].dataType) + bld.Reserve(sizeHint) + g.builders = append(g.builders, bld) + g.appenders = append(g.appenders, ppeAppender(bld, col.slot, e.slots[col.slot])) + } + return g +} + +// ppeAppender binds the builder type once, so the per-row hot path does not +// type switch. +func ppeAppender(bld array.Builder, slot int, s ppeSlot) func(r *ppeRow) { + if s.repeated { + lb := bld.(*array.ListBuilder) + elem := ppeAppender(lb.ValueBuilder(), 0, ppeSlot{kind: s.kind}) + scratch := &ppeRow{bits: make([]uint64, 1)} + return func(r *ppeRow) { + lb.Append(true) + for _, v := range r.lists[slot] { + scratch.bits[0] = v + elem(scratch) + } + } + } + switch b := bld.(type) { + case *array.Int64Builder: + switch s.kind { + case protoreflect.Sint64Kind: + return func(r *ppeRow) { b.Append(protowire.DecodeZigZag(r.bits[slot])) } + default: + return func(r *ppeRow) { b.Append(int64(r.bits[slot])) } + } + case *array.Int32Builder: + switch s.kind { + case protoreflect.Sint32Kind: + // Inlined rather than routed through ppeSignedInt: this runs per row, + // and the boxing the shared helper used to do allocated on every one. + return func(r *ppeRow) { b.Append(int32(protowire.DecodeZigZag(r.bits[slot] & math.MaxUint32))) } + case protoreflect.Sfixed32Kind: + return func(r *ppeRow) { b.Append(int32(uint32(r.bits[slot]))) } + default: + return func(r *ppeRow) { b.Append(int32(r.bits[slot])) } + } + case *array.Uint64Builder: + return func(r *ppeRow) { b.Append(r.bits[slot]) } + case *array.Uint32Builder: + return func(r *ppeRow) { b.Append(uint32(r.bits[slot])) } + case *array.Float64Builder: + return func(r *ppeRow) { b.Append(math.Float64frombits(r.bits[slot])) } + case *array.Float32Builder: + return func(r *ppeRow) { b.Append(math.Float32frombits(uint32(r.bits[slot]))) } + case *array.BooleanBuilder: + return func(r *ppeRow) { b.Append(r.bits[slot] != 0) } + case *array.StringBuilder: + return func(r *ppeRow) { b.BinaryBuilder.Append(r.bytes[slot]) } + case *array.BinaryBuilder: + return func(r *ppeRow) { b.Append(r.bytes[slot]) } + } + panic(fmt.Sprintf("unsupported builder %T", bld)) +} + +// ppeSignedInt is the one place raw slot bits become a signed integer, shared by +// the partition key and the Int32/Int64 appenders so the directory a row lands in +// can never contradict the row's own timestamp column. +func ppeSignedInt(kind protoreflect.Kind, bits uint64) int64 { + switch kind { + case protoreflect.Sint64Kind: + return protowire.DecodeZigZag(bits) + case protoreflect.Sint32Kind: + return int64(int32(protowire.DecodeZigZag(bits & math.MaxUint32))) + case protoreflect.Sfixed32Kind: + // Fixed32 is stored zero-extended; the signed form is the low 32 bits. + return int64(int32(uint32(bits))) + } + return int64(bits) +} + +func (e *protobufParquetEncoder) partitionKey(r *ppeRow, s *ppeSlot, lastSec *int64, lastKey *string) string { + v := ppeSignedInt(s.kind, r.bits[e.partitionSlot]) + sec := v / e.partitionDivisor + if v < 0 && v%e.partitionDivisor != 0 { + sec-- + } + if sec != *lastSec || *lastKey == "" { + *lastSec = sec + *lastKey = time.Unix(sec, 0).UTC().Format(e.partitionLayout) + } + return *lastKey +} + +func (e *protobufParquetEncoder) ProcessBatch(ctx context.Context, batch service.MessageBatch) ([]service.MessageBatch, error) { + if len(batch) == 0 { + return nil, nil + } + mem := memory.NewGoAllocator() + row := &ppeRow{bits: make([]uint64, len(e.slots)), bytes: make([][]byte, len(e.slots)), lists: make([][]uint64, len(e.slots))} + + var groups []*ppeGroup + byKey := map[string]*ppeGroup{} + var current *ppeGroup + lastSec, lastKey := int64(0), "" + invalid := 0 + + processed := 0 + for _, msg := range batch { + raw, err := msg.AsBytes() + if err == nil { + err = e.decode(raw, row) + } + if err != nil { + if !e.skipInvalid { + return nil, fmt.Errorf("failed to decode protobuf message: %w", err) + } + invalid++ + continue + } + key := "" + if e.partitionSlot >= 0 { + key = e.partitionKey(row, &e.slots[e.partitionSlot], &lastSec, &lastKey) + } + if current == nil || current.key != key { + if current = byKey[key]; current == nil { + // Only the rows still to come can land in a group created now; + // hinting the whole batch for every key multiplies peak memory by + // the number of partitions a batch straddles. + current = e.newGroup(msg, key, mem, len(batch)-processed) + byKey[key] = current + groups = append(groups, current) + } + } + for _, appendCol := range current.appenders { + appendCol(row) + } + current.rows++ + processed++ + } + if invalid > 0 { + e.mInvalid.Incr(int64(invalid)) + e.log.Errorf("Dropped %d of %d messages that could not be decoded as protobuf", invalid, len(batch)) + } + + out := make(service.MessageBatch, 0, len(groups)) + for _, g := range groups { + data, err := e.write(g) + if err != nil { + return nil, err + } + m := g.first.Copy() + m.SetBytes(data) + if e.partitionSlot >= 0 { + m.MetaSetMut(e.partitionMetaKey, g.key) + } + out = append(out, m) + } + if len(out) == 0 { + return nil, nil + } + return []service.MessageBatch{out}, nil +} + +func (e *protobufParquetEncoder) write(g *ppeGroup) ([]byte, error) { + cols := make([]arrow.Array, len(g.builders)) + for i, b := range g.builders { + cols[i] = b.NewArray() + b.Release() + } + rec := array.NewRecordBatch(e.schema, cols, int64(g.rows)) + defer rec.Release() + for _, c := range cols { + c.Release() + } + + var buf bytes.Buffer + w, err := pqarrow.NewFileWriter(e.schema, &buf, e.props, pqarrow.DefaultWriterProps()) + if err != nil { + return nil, fmt.Errorf("failed to create parquet writer: %w", err) + } + if err := w.Write(rec); err != nil { + _ = w.Close() + return nil, fmt.Errorf("failed to write parquet row group: %w", err) + } + if err := w.Close(); err != nil { + return nil, fmt.Errorf("failed to finalize parquet file: %w", err) + } + return buf.Bytes(), nil +} + +func (e *protobufParquetEncoder) Close(ctx context.Context) error { + return nil +} diff --git a/internal/impl/protobuf/processor_protobuf_parquet_test.go b/internal/impl/protobuf/processor_protobuf_parquet_test.go new file mode 100644 index 0000000000..4cf59da88c --- /dev/null +++ b/internal/impl/protobuf/processor_protobuf_parquet_test.go @@ -0,0 +1,825 @@ +package protobuf + +import ( + "bytes" + "context" + "fmt" + "math" + "os" + "path/filepath" + "testing" + "time" + + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/memory" + "github.com/apache/arrow-go/v18/parquet/compress" + "github.com/apache/arrow-go/v18/parquet/file" + "github.com/apache/arrow-go/v18/parquet/pqarrow" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "google.golang.org/protobuf/encoding/protowire" + "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/reflect/protoreflect" + "google.golang.org/protobuf/types/dynamicpb" + + "github.com/warpstreamlabs/bento/public/service" +) + +const ppeTestProto = `syntax = "proto3"; +package bento.test; + +message Tick { + int64 time_us = 1; + int64 exchange = 2; + string symbol = 3; + double price = 4; + int64 local_timestamp_us = 5; + bool is_snapshot = 6; + repeated double levels = 7; + sint64 delta = 8; + int32 venue = 9; +} +` + +func ppeTestDir(t testing.TB) string { + t.Helper() + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "tick.proto"), []byte(ppeTestProto), 0o644)) + return dir +} + +func ppeTickType(t testing.TB, dir string) protoreflect.MessageType { + t.Helper() + files, _, err := loadDescriptors(service.MockResources().FS(), []string{dir}) + require.NoError(t, err) + d, err := files.FindDescriptorByName("bento.test.Tick") + require.NoError(t, err) + return dynamicpb.NewMessageType(d.(protoreflect.MessageDescriptor)) +} + +type ppeTick struct { + timeUs, exchange, localTs, delta int64 + symbol string + price float64 + snapshot bool + levels []float64 + venue int32 +} + +func ppeMarshal(t testing.TB, mt protoreflect.MessageType, v ppeTick) []byte { + t.Helper() + m := mt.New() + fields := m.Descriptor().Fields() + set := func(name string, val protoreflect.Value) { m.Set(fields.ByName(protoreflect.Name(name)), val) } + set("time_us", protoreflect.ValueOfInt64(v.timeUs)) + set("exchange", protoreflect.ValueOfInt64(v.exchange)) + set("symbol", protoreflect.ValueOfString(v.symbol)) + set("price", protoreflect.ValueOfFloat64(v.price)) + set("local_timestamp_us", protoreflect.ValueOfInt64(v.localTs)) + set("is_snapshot", protoreflect.ValueOfBool(v.snapshot)) + set("delta", protoreflect.ValueOfInt64(v.delta)) + set("venue", protoreflect.ValueOfInt32(v.venue)) + if len(v.levels) > 0 { + l := m.Mutable(fields.ByName("levels")).List() + for _, x := range v.levels { + l.Append(protoreflect.ValueOfFloat64(x)) + } + } + b, err := proto.Marshal(m.Interface()) + require.NoError(t, err) + return b +} + +func ppeNewProc(t testing.TB, dir, extra string) *protobufParquetEncoder { + t.Helper() + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf(` +message: bento.test.Tick +import_paths: [ %v ] +columns: + - { name: exchange, constant: binance-futures } + - { name: symbol } + - { name: time, field: time_us } + - { name: local_timestamp, field: local_timestamp_us } + - { name: price } + - { name: is_snapshot } + - { name: levels } + - { name: delta } + - { name: venue } +%v`, dir, extra), nil) + require.NoError(t, err) + p, err := newProtobufParquetEncoder(conf, service.MockResources()) + require.NoError(t, err) + return p +} + +func ppeReadTable(t testing.TB, data []byte) arrow.Table { + t.Helper() + rdr, err := file.NewParquetReader(bytes.NewReader(data)) + require.NoError(t, err) + fr, err := pqarrow.NewFileReader(rdr, pqarrow.ArrowReadProperties{}, memory.DefaultAllocator) + require.NoError(t, err) + tbl, err := fr.ReadTable(context.Background()) + require.NoError(t, err) + t.Cleanup(tbl.Release) + assert.Equal(t, 1, rdr.NumRowGroups(), "a batch must become a single row group") + return tbl +} + +func ppeCol(t testing.TB, tbl arrow.Table, name string) arrow.Array { + t.Helper() + idx := tbl.Schema().FieldIndices(name) + require.Len(t, idx, 1, "column %v", name) + require.Len(t, tbl.Column(idx[0]).Data().Chunks(), 1) + return tbl.Column(idx[0]).Data().Chunk(0) +} + +func TestProtobufParquetEncodeRoundTrip(t *testing.T) { + dir := ppeTestDir(t) + mt := ppeTickType(t, dir) + p := ppeNewProc(t, dir, "") + + ticks := []ppeTick{ + {timeUs: 1_789_502_374_870_000, exchange: 11, symbol: "BTC-USDT", price: 101.5, localTs: 1_789_502_374_995_869, snapshot: true, levels: []float64{1.5, 2.5, 3.5}, delta: -42, venue: -7}, + // Zero values are not on the wire in proto3; they must still come back as zeroes, not nulls. + {symbol: "ETH-USDT"}, + {timeUs: 3, exchange: 11, symbol: "BTC-USDT", price: -0.25, localTs: 4, levels: []float64{9}}, + } + batch := service.MessageBatch{} + for i, tk := range ticks { + m := service.NewMessage(ppeMarshal(t, mt, tk)) + m.MetaSetMut("kafka_offset", fmt.Sprint(i)) + batch = append(batch, m) + } + + out, err := p.ProcessBatch(context.Background(), batch) + require.NoError(t, err) + require.Len(t, out, 1) + require.Len(t, out[0], 1) + + // The output keeps the first message's metadata, like parquet_encode. + off, ok := out[0][0].MetaGetMut("kafka_offset") + require.True(t, ok) + assert.Equal(t, "0", off) + + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + require.EqualValues(t, len(ticks), tbl.NumRows()) + + var names []string + for _, f := range tbl.Schema().Fields() { + names = append(names, f.Name) + } + assert.Equal(t, []string{"exchange", "symbol", "time", "local_timestamp", "price", "is_snapshot", "levels", "delta", "venue"}, names) + + exch := ppeCol(t, tbl, "exchange").(*array.String) + sym := ppeCol(t, tbl, "symbol").(*array.String) + tm := ppeCol(t, tbl, "time").(*array.Int64) + lts := ppeCol(t, tbl, "local_timestamp").(*array.Int64) + price := ppeCol(t, tbl, "price").(*array.Float64) + snap := ppeCol(t, tbl, "is_snapshot").(*array.Boolean) + levels := ppeCol(t, tbl, "levels").(*array.List) + delta := ppeCol(t, tbl, "delta").(*array.Int64) + venue := ppeCol(t, tbl, "venue").(*array.Int32) + for i, tk := range ticks { + assert.Equal(t, "binance-futures", exch.Value(i)) + assert.Equal(t, tk.symbol, sym.Value(i)) + assert.Equal(t, tk.timeUs, tm.Value(i)) + assert.Equal(t, tk.localTs, lts.Value(i)) + assert.Equal(t, tk.price, price.Value(i)) + assert.Equal(t, tk.snapshot, snap.Value(i)) + assert.Equal(t, tk.delta, delta.Value(i)) + assert.Equal(t, tk.venue, venue.Value(i)) + assert.False(t, levels.IsNull(i)) + start, end := levels.ValueOffsets(i) + var got []float64 + for j := start; j < end; j++ { + got = append(got, levels.ListValues().(*array.Float64).Value(int(j))) + } + if len(tk.levels) == 0 { + assert.Empty(t, got) + } else { + assert.Equal(t, tk.levels, got) + } + } +} + +func TestProtobufParquetEncodePartitionsByDay(t *testing.T) { + dir := ppeTestDir(t) + mt := ppeTickType(t, dir) + p := ppeNewProc(t, dir, ` +partition: + field: local_timestamp_us + unit: us + layout: year=2006/month=01/day=02 + metadata_key: day +`) + day1 := time.Date(2026, 9, 15, 23, 59, 59, 0, time.UTC).UnixMicro() + day2 := time.Date(2026, 9, 16, 0, 0, 0, 0, time.UTC).UnixMicro() + var batch service.MessageBatch + for _, ts := range []int64{day1, day2, day1 - 1, day2 + 5, day2 + 6} { + batch = append(batch, service.NewMessage(ppeMarshal(t, mt, ppeTick{symbol: "X", localTs: ts}))) + } + out, err := p.ProcessBatch(context.Background(), batch) + require.NoError(t, err) + require.Len(t, out, 1) + require.Len(t, out[0], 2, "one file per day") + + wantKeys := []string{"year=2026/month=09/day=15", "year=2026/month=09/day=16"} + wantRows := []int64{2, 3} + for i, m := range out[0] { + key, ok := m.MetaGetMut("day") + require.True(t, ok) + assert.Equal(t, wantKeys[i], key) + data, err := m.AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + assert.Equal(t, wantRows[i], tbl.NumRows()) + lts := ppeCol(t, tbl, "local_timestamp").(*array.Int64) + for j := 0; j < lts.Len(); j++ { + assert.Equal(t, wantKeys[i], time.UnixMicro(lts.Value(j)).UTC().Format("year=2006/month=01/day=02"), "row landed in the wrong partition") + } + } +} + +func TestProtobufParquetEncodeInvalidMessages(t *testing.T) { + dir := ppeTestDir(t) + mt := ppeTickType(t, dir) + good := ppeMarshal(t, mt, ppeTick{symbol: "OK", timeUs: 1}) + truncated := good[:len(good)-1] + + p := ppeNewProc(t, dir, "") + out, err := p.ProcessBatch(context.Background(), service.MessageBatch{ + service.NewMessage(truncated), service.NewMessage(good), + }) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + assert.EqualValues(t, 1, tbl.NumRows(), "the undecodable message is dropped, the valid one is kept") + + strict := ppeNewProc(t, dir, "skip_invalid_messages: false") + _, err = strict.ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(truncated)}) + require.Error(t, err) + + out, err = p.ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(truncated)}) + require.NoError(t, err) + assert.Empty(t, out, "a batch with no decodable message produces no file") +} + +func TestProtobufParquetEncodeConfigErrors(t *testing.T) { + dir := ppeTestDir(t) + for name, conf := range map[string]string{ + "unknown field": "columns: [ { name: nope } ]", + "field and constant": "columns: [ { name: s, field: symbol, constant: x } ]", + "string partition field": "columns: [ { name: symbol } ]\npartition: { field: symbol, layout: '2006' }", + "unknown message": "columns: [ { name: symbol } ]", + } { + t.Run(name, func(t *testing.T) { + msg := "bento.test.Tick" + if name == "unknown message" { + msg = "bento.test.Nope" + } + parsed, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf("message: %v\nimport_paths: [ %v ]\n%v", msg, dir, conf), nil) + require.NoError(t, err) + _, err = newProtobufParquetEncoder(parsed, service.MockResources()) + require.Error(t, err) + }) + } +} + +func BenchmarkProtobufParquetEncode(b *testing.B) { + dir := ppeTestDir(b) + mt := ppeTickType(b, dir) + p := ppeNewProc(b, dir, "partition: { field: local_timestamp_us, layout: 'year=2006/month=01/day=02' }") + base := time.Date(2026, 9, 15, 12, 0, 0, 0, time.UTC).UnixMicro() + const n = 50_000 + batch := make(service.MessageBatch, n) + for i := range batch { + batch[i] = service.NewMessage(ppeMarshal(b, mt, ppeTick{ + timeUs: base + int64(i)*1000, exchange: 11, symbol: fmt.Sprintf("SYM%d-USDT", i%400), + price: 100 + float64(i%1000)/7, localTs: base + int64(i)*1000 + 250, + })) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := p.ProcessBatch(context.Background(), batch); err != nil { + b.Fatal(err) + } + } + b.ReportMetric(float64(n*b.N)/b.Elapsed().Seconds(), "msgs/s") +} + +// ------------------------------------------------------------------------------ +// Wire-format coverage. Every scalar kind the encoder claims to support gets its +// own decode branch and its own Arrow builder, and several of them reinterpret +// bits (zigzag, sign-extension, float bit patterns) in ways a happy-path fixture +// of int64/double/string never touches. + +const ppeKindsProto = `syntax = "proto3"; +package bento.test; + +enum Venue { + VENUE_UNSPECIFIED = 0; + VENUE_SPOT = 1; + VENUE_PERP = 2; +} + +message Kinds { + int32 f_int32 = 1; + int64 f_int64 = 2; + uint32 f_uint32 = 3; + uint64 f_uint64 = 4; + sint32 f_sint32 = 5; + sint64 f_sint64 = 6; + fixed32 f_fixed32 = 7; + fixed64 f_fixed64 = 8; + sfixed32 f_sfixed32 = 9; + sfixed64 f_sfixed64 = 10; + float f_float = 11; + double f_double = 12; + bool f_bool = 13; + string f_string = 14; + bytes f_bytes = 15; + // Field 16 and up need a two-byte tag, which is its own decode path. + Venue f_enum = 16; + repeated int32 r_int32 = 17; + repeated float r_float = 18; + repeated string r_string = 19; + map m_counts = 20; +} +` + +func ppeKindsDir(t testing.TB) string { + t.Helper() + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "kinds.proto"), []byte(ppeKindsProto), 0o644)) + return dir +} + +func ppeKindsType(t testing.TB, dir string) protoreflect.MessageType { + t.Helper() + files, _, err := loadDescriptors(service.MockResources().FS(), []string{dir}) + require.NoError(t, err) + d, err := files.FindDescriptorByName("bento.test.Kinds") + require.NoError(t, err) + return dynamicpb.NewMessageType(d.(protoreflect.MessageDescriptor)) +} + +func ppeKindsProc(t testing.TB, dir string) *protobufParquetEncoder { + t.Helper() + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf(` +message: bento.test.Kinds +import_paths: [ %v ] +columns: + - { name: f_int32 } + - { name: f_int64 } + - { name: f_uint32 } + - { name: f_uint64 } + - { name: f_sint32 } + - { name: f_sint64 } + - { name: f_fixed32 } + - { name: f_fixed64 } + - { name: f_sfixed32 } + - { name: f_sfixed64 } + - { name: f_float } + - { name: f_double } + - { name: f_bool } + - { name: f_string } + - { name: f_bytes } + - { name: f_enum } + - { name: r_int32 } + - { name: r_float } +`, dir), nil) + require.NoError(t, err) + p, err := newProtobufParquetEncoder(conf, service.MockResources()) + require.NoError(t, err) + return p +} + +// Extremes on purpose: the min/max of each width is where a missing +// sign-extension or a stray uint64->int64 conversion shows up. +func TestProtobufParquetEncodeEveryScalarKind(t *testing.T) { + dir := ppeKindsDir(t) + mt := ppeKindsType(t, dir) + + m := mt.New() + fields := m.Descriptor().Fields() + set := func(name string, v protoreflect.Value) { m.Set(fields.ByName(protoreflect.Name(name)), v) } + set("f_int32", protoreflect.ValueOfInt32(math.MinInt32)) + set("f_int64", protoreflect.ValueOfInt64(math.MinInt64)) + set("f_uint32", protoreflect.ValueOfUint32(math.MaxUint32)) + set("f_uint64", protoreflect.ValueOfUint64(math.MaxUint64)) + set("f_sint32", protoreflect.ValueOfInt32(math.MinInt32)) + set("f_sint64", protoreflect.ValueOfInt64(math.MinInt64)) + set("f_fixed32", protoreflect.ValueOfUint32(math.MaxUint32)) + set("f_fixed64", protoreflect.ValueOfUint64(math.MaxUint64)) + set("f_sfixed32", protoreflect.ValueOfInt32(math.MinInt32)) + set("f_sfixed64", protoreflect.ValueOfInt64(math.MinInt64)) + set("f_float", protoreflect.ValueOfFloat32(-1.5)) + set("f_double", protoreflect.ValueOfFloat64(-math.MaxFloat64)) + set("f_bool", protoreflect.ValueOfBool(true)) + set("f_string", protoreflect.ValueOfString("héllo")) + set("f_bytes", protoreflect.ValueOfBytes([]byte{0x00, 0x01, 0xff})) + set("f_enum", protoreflect.ValueOfEnum(2)) + ints := m.Mutable(fields.ByName("r_int32")).List() + for _, v := range []int32{1, -2, math.MaxInt32} { + ints.Append(protoreflect.ValueOfInt32(v)) + } + floats := m.Mutable(fields.ByName("r_float")).List() + for _, v := range []float32{1.5, -2.5} { + floats.Append(protoreflect.ValueOfFloat32(v)) + } + raw, err := proto.Marshal(m.Interface()) + require.NoError(t, err) + + out, err := ppeKindsProc(t, dir).ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + require.Len(t, out, 1) + require.Len(t, out[0], 1) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + require.EqualValues(t, 1, tbl.NumRows()) + + assert.Equal(t, int32(math.MinInt32), ppeCol(t, tbl, "f_int32").(*array.Int32).Value(0)) + assert.Equal(t, int64(math.MinInt64), ppeCol(t, tbl, "f_int64").(*array.Int64).Value(0)) + assert.Equal(t, uint32(math.MaxUint32), ppeCol(t, tbl, "f_uint32").(*array.Uint32).Value(0)) + assert.Equal(t, uint64(math.MaxUint64), ppeCol(t, tbl, "f_uint64").(*array.Uint64).Value(0)) + assert.Equal(t, int32(math.MinInt32), ppeCol(t, tbl, "f_sint32").(*array.Int32).Value(0), "sint32 must be zigzag-decoded") + assert.Equal(t, int64(math.MinInt64), ppeCol(t, tbl, "f_sint64").(*array.Int64).Value(0), "sint64 must be zigzag-decoded") + assert.Equal(t, uint32(math.MaxUint32), ppeCol(t, tbl, "f_fixed32").(*array.Uint32).Value(0)) + assert.Equal(t, uint64(math.MaxUint64), ppeCol(t, tbl, "f_fixed64").(*array.Uint64).Value(0)) + assert.Equal(t, int32(math.MinInt32), ppeCol(t, tbl, "f_sfixed32").(*array.Int32).Value(0), "sfixed32 keeps its sign") + assert.Equal(t, int64(math.MinInt64), ppeCol(t, tbl, "f_sfixed64").(*array.Int64).Value(0), "sfixed64 keeps its sign") + assert.Equal(t, float32(-1.5), ppeCol(t, tbl, "f_float").(*array.Float32).Value(0)) + assert.Equal(t, -math.MaxFloat64, ppeCol(t, tbl, "f_double").(*array.Float64).Value(0)) + assert.True(t, ppeCol(t, tbl, "f_bool").(*array.Boolean).Value(0)) + assert.Equal(t, "héllo", ppeCol(t, tbl, "f_string").(*array.String).Value(0)) + assert.Equal(t, []byte{0x00, 0x01, 0xff}, ppeCol(t, tbl, "f_bytes").(*array.Binary).Value(0)) + assert.Equal(t, int32(2), ppeCol(t, tbl, "f_enum").(*array.Int32).Value(0), "an enum lands as its number") + + rInts := ppeCol(t, tbl, "r_int32").(*array.List) + start, end := rInts.ValueOffsets(0) + var gotInts []int32 + for i := start; i < end; i++ { + gotInts = append(gotInts, rInts.ListValues().(*array.Int32).Value(int(i))) + } + assert.Equal(t, []int32{1, -2, math.MaxInt32}, gotInts) + + rFloats := ppeCol(t, tbl, "r_float").(*array.List) + start, end = rFloats.ValueOffsets(0) + var gotFloats []float32 + for i := start; i < end; i++ { + gotFloats = append(gotFloats, rFloats.ListValues().(*array.Float32).Value(int(i))) + } + assert.Equal(t, []float32{1.5, -2.5}, gotFloats) +} + +// proto3 packs repeated scalars, but nothing requires a producer to: proto2 +// encoders, hand-rolled writers and some language runtimes emit one tag per +// element, and a decoder that only understands the packed form silently returns +// an empty list for them. +func TestProtobufParquetEncodeUnpackedRepeatedFields(t *testing.T) { + dir := ppeKindsDir(t) + + var raw []byte + for _, v := range []int32{7, -8, 9} { + raw = protowire.AppendTag(raw, 17, protowire.VarintType) + raw = protowire.AppendVarint(raw, uint64(uint32(v))) + } + for _, v := range []float32{0.5, -0.25} { + raw = protowire.AppendTag(raw, 18, protowire.Fixed32Type) + raw = protowire.AppendFixed32(raw, math.Float32bits(v)) + } + + out, err := ppeKindsProc(t, dir).ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + + ints := ppeCol(t, tbl, "r_int32").(*array.List) + start, end := ints.ValueOffsets(0) + var gotInts []int32 + for i := start; i < end; i++ { + gotInts = append(gotInts, ints.ListValues().(*array.Int32).Value(int(i))) + } + assert.Equal(t, []int32{7, -8, 9}, gotInts, "unpacked varint elements must all be kept") + + floats := ppeCol(t, tbl, "r_float").(*array.List) + start, end = floats.ValueOffsets(0) + var gotFloats []float32 + for i := start; i < end; i++ { + gotFloats = append(gotFloats, floats.ListValues().(*array.Float32).Value(int(i))) + } + assert.Equal(t, []float32{0.5, -0.25}, gotFloats, "unpacked fixed32 elements must all be kept") +} + +// A record written by a newer producer carries fields this config never mapped. +// Skipping them has to consume exactly the right number of bytes for each wire +// type, or every field after the unknown one decodes as garbage. +func TestProtobufParquetEncodeSkipsUnknownFields(t *testing.T) { + dir := ppeKindsDir(t) + + var raw []byte + raw = protowire.AppendTag(raw, 900, protowire.VarintType) + raw = protowire.AppendVarint(raw, math.MaxUint64) + raw = protowire.AppendTag(raw, 901, protowire.Fixed64Type) + raw = protowire.AppendFixed64(raw, 0xdeadbeefcafef00d) + raw = protowire.AppendTag(raw, 902, protowire.Fixed32Type) + raw = protowire.AppendFixed32(raw, 0xfeedface) + raw = protowire.AppendTag(raw, 903, protowire.BytesType) + raw = protowire.AppendBytes(raw, []byte("an unmapped string")) + // Only now the field the config actually asks for. + raw = protowire.AppendTag(raw, 2, protowire.VarintType) + raw = protowire.AppendVarint(raw, 4242) + + out, err := ppeKindsProc(t, dir).ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + require.EqualValues(t, 1, tbl.NumRows()) + assert.Equal(t, int64(4242), ppeCol(t, tbl, "f_int64").(*array.Int64).Value(0), + "the mapped field after four unknown ones must still decode") +} + +// A field repeated on the wire for a singular column is legal protobuf and means +// "last one wins" -- merge semantics, not an error. +func TestProtobufParquetEncodeLastValueWinsForSingularFields(t *testing.T) { + dir := ppeKindsDir(t) + + var raw []byte + for _, v := range []uint64{1, 2, 3} { + raw = protowire.AppendTag(raw, 2, protowire.VarintType) + raw = protowire.AppendVarint(raw, v) + } + + out, err := ppeKindsProc(t, dir).ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + assert.Equal(t, int64(3), ppeCol(t, tbl, "f_int64").(*array.Int64).Value(0)) +} + +// The four supported units have to divide the same instant to the same day, or a +// pipeline configured in ms silently files its rows under the wrong date. +func TestProtobufParquetEncodePartitionUnits(t *testing.T) { + dir := ppeTestDir(t) + mt := ppeTickType(t, dir) + instant := time.Date(2026, 9, 15, 13, 45, 30, 500_000_000, time.UTC) + + for unit, ts := range map[string]int64{ + "s": instant.Unix(), + "ms": instant.UnixMilli(), + "us": instant.UnixMicro(), + "ns": instant.UnixNano(), + } { + t.Run(unit, func(t *testing.T) { + p := ppeNewProc(t, dir, fmt.Sprintf(` +partition: + field: local_timestamp_us + unit: %v + layout: year=2006/month=01/day=02 +`, unit)) + out, err := p.ProcessBatch(context.Background(), service.MessageBatch{ + service.NewMessage(ppeMarshal(t, mt, ppeTick{symbol: "X", localTs: ts})), + }) + require.NoError(t, err) + require.Len(t, out[0], 1) + key, ok := out[0][0].MetaGetMut("partition") + require.True(t, ok) + assert.Equal(t, "year=2026/month=09/day=15", key) + }) + } +} + +// The codec is not observable from the decoded rows -- a file written with the +// wrong one reads back identically -- so it has to be asserted against the +// Parquet metadata, which is also what a downstream reader negotiates on. +func TestProtobufParquetEncodeCompressionOptions(t *testing.T) { + dir := ppeTestDir(t) + mt := ppeTickType(t, dir) + + for _, tc := range []struct { + name string + extra string + want compress.Compression + }{ + {name: "default is zstd", want: compress.Codecs.Zstd}, + {name: "snappy", extra: "compression: snappy", want: compress.Codecs.Snappy}, + {name: "uncompressed", extra: "compression: uncompressed", want: compress.Codecs.Uncompressed}, + {name: "gzip with a level", extra: "compression: gzip\ncompression_level: 1", want: compress.Codecs.Gzip}, + {name: "zstd without dictionary encoding", extra: "dictionary: false", want: compress.Codecs.Zstd}, + } { + t.Run(tc.name, func(t *testing.T) { + p := ppeNewProc(t, dir, tc.extra) + out, err := p.ProcessBatch(context.Background(), service.MessageBatch{ + service.NewMessage(ppeMarshal(t, mt, ppeTick{symbol: "X", price: 1.25})), + }) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + + rdr, err := file.NewParquetReader(bytes.NewReader(data)) + require.NoError(t, err) + t.Cleanup(func() { _ = rdr.Close() }) + chunk, err := rdr.MetaData().RowGroup(0).ColumnChunk(0) + require.NoError(t, err) + assert.Equal(t, tc.want, chunk.Compression()) + + // Whatever the codec, the rows still have to survive it. + assert.Equal(t, "X", ppeCol(t, ppeReadTable(t, data), "symbol").(*array.String).Value(0)) + }) + } +} + +// These two are rejected at construction rather than mishandled at runtime, and +// the repeated-string guard is load-bearing: a repeated string arrives as one +// length-delimited chunk per element, which the packed-scalar path would happily +// misread as a run of varints. +func TestProtobufParquetEncodeRejectsUnsupportedFieldShapes(t *testing.T) { + dir := ppeKindsDir(t) + for name, column := range map[string]string{ + "repeated string": "r_string", + "map": "m_counts", + } { + t.Run(name, func(t *testing.T) { + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf( + "message: bento.test.Kinds\nimport_paths: [ %v ]\ncolumns: [ { name: %v } ]", dir, column), nil) + require.NoError(t, err) + _, err = newProtobufParquetEncoder(conf, service.MockResources()) + require.Error(t, err) + }) + } +} + +// NewStringEnumField only attaches a lint rule, so an out-of-enum value reaches +// the constructor intact whenever linting is skipped (the streams API, --chilled, +// a programmatic ParseYAML). Both of these used to be unchecked map lookups that +// returned the zero value: the unit became divisor 0 and panicked with an integer +// divide by zero on the first message, and the codec became "uncompressed" and +// silently dropped compression on every file written. +func TestProtobufParquetEncodeRejectsOutOfEnumOptions(t *testing.T) { + dir := ppeTestDir(t) + for name, extra := range map[string]string{ + "unknown partition unit": "partition: { field: local_timestamp_us, unit: seconds, layout: '2006' }", + "unknown compression": "compression: zippy", + } { + t.Run(name, func(t *testing.T) { + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf( + "message: bento.test.Tick\nimport_paths: [ %v ]\ncolumns: [ { name: symbol } ]\n%v", dir, extra), nil) + require.NoError(t, err) + _, err = newProtobufParquetEncoder(conf, service.MockResources()) + require.Error(t, err) + }) + } +} + +// The partition key and the column value are derived from the same bits by two +// different code paths, and they have to agree: a row whose own timestamp column +// says 1969 must not be filed under 2106. sfixed32 is where they diverged -- +// the column appender reinterpreted the low 32 bits as signed, the partition key +// did not. +func TestProtobufParquetEncodePartitionKeyMatchesColumnForSignedKinds(t *testing.T) { + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "signed.proto"), []byte(`syntax = "proto3"; +package bento.test; +message Signed { + sfixed32 ts_sfixed32 = 1; + sint64 ts_sint64 = 2; + sint32 ts_sint32 = 3; + int32 ts_int32 = 4; +} +`), 0o644)) + files, _, err := loadDescriptors(service.MockResources().FS(), []string{dir}) + require.NoError(t, err) + d, err := files.FindDescriptorByName("bento.test.Signed") + require.NoError(t, err) + mt := dynamicpb.NewMessageType(d.(protoreflect.MessageDescriptor)) + + // One day before the epoch: every signed kind must land on 1969-12-31. + const beforeEpoch = int64(-86400) + for _, field := range []string{"ts_sfixed32", "ts_sint64", "ts_sint32", "ts_int32"} { + t.Run(field, func(t *testing.T) { + m := mt.New() + fd := m.Descriptor().Fields().ByName(protoreflect.Name(field)) + if fd.Kind() == protoreflect.Sint64Kind { + m.Set(fd, protoreflect.ValueOfInt64(beforeEpoch)) + } else { + m.Set(fd, protoreflect.ValueOfInt32(int32(beforeEpoch))) + } + raw, err := proto.Marshal(m.Interface()) + require.NoError(t, err) + + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf(` +message: bento.test.Signed +import_paths: [ %v ] +columns: [ { name: ts, field: %v } ] +partition: { field: %v, unit: s, layout: 'year=2006/month=01/day=02' } +`, dir, field, field), nil) + require.NoError(t, err) + p, err := newProtobufParquetEncoder(conf, service.MockResources()) + require.NoError(t, err) + + out, err := p.ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + require.Len(t, out[0], 1) + key, ok := out[0][0].MetaGetMut("partition") + require.True(t, ok) + assert.Equal(t, "year=1969/month=12/day=31", key, "partition key disagrees with the row's own timestamp") + + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + col := ppeCol(t, tbl, "ts") + var got int64 + switch c := col.(type) { + case *array.Int64: + got = c.Value(0) + case *array.Int32: + got = int64(c.Value(0)) + default: + t.Fatalf("unexpected column type %T", col) + } + assert.Equal(t, beforeEpoch, got) + }) + } +} + +// A length-delimited payload arriving on a packable repeated field is the packed +// encoding -- the wire format offers no way to tell it apart from a string sent +// on the same field number. This pins that the decoder agrees with the protobuf +// runtime rather than inventing a stricter rule of its own. +func TestProtobufParquetEncodePackedDecodeMatchesProtoUnmarshal(t *testing.T) { + dir := ppeKindsDir(t) + mt := ppeKindsType(t, dir) + + var raw []byte + raw = protowire.AppendTag(raw, 17, protowire.BytesType) + raw = protowire.AppendBytes(raw, []byte("BTCUSDT")) + + reference := mt.New().Interface() + require.NoError(t, proto.Unmarshal(raw, reference)) + refList := reference.ProtoReflect().Get(mt.Descriptor().Fields().ByName("r_int32")).List() + var want []int32 + for i := 0; i < refList.Len(); i++ { + want = append(want, int32(refList.Get(i).Int())) + } + require.NotEmpty(t, want, "the protobuf runtime itself reads these bytes as packed varints") + + out, err := ppeKindsProc(t, dir).ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + list := ppeCol(t, tbl, "r_int32").(*array.List) + start, end := list.ValueOffsets(0) + var got []int32 + for i := start; i < end; i++ { + got = append(got, list.ListValues().(*array.Int32).Value(int(i))) + } + assert.Equal(t, want, got) +} + +// Field numbers go up to 536,870,911, and a slot table indexed by field number +// allocates one entry per number up to the highest one used -- 8MB for a single +// column on field 1,000,000, and gigabytes near the top of the range. +func TestProtobufParquetEncodeHandlesSparseFieldNumbers(t *testing.T) { + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "sparse.proto"), []byte(`syntax = "proto3"; +package bento.test; +message Sparse { + string symbol = 1; + int64 ts = 536870911; +} +`), 0o644)) + conf, err := protobufParquetEncodeSpec().ParseYAML(fmt.Sprintf(` +message: bento.test.Sparse +import_paths: [ %v ] +columns: [ { name: symbol }, { name: ts } ] +`, dir), nil) + require.NoError(t, err) + + p, err := newProtobufParquetEncoder(conf, service.MockResources()) + require.NoError(t, err) + assert.Nil(t, p.slotByField, "a schema numbered this high must not build a dense slot table") + + var raw []byte + raw = protowire.AppendTag(raw, 1, protowire.BytesType) + raw = protowire.AppendBytes(raw, []byte("BTC-USDT")) + raw = protowire.AppendTag(raw, 536870911, protowire.VarintType) + raw = protowire.AppendVarint(raw, 1234) + + out, err := p.ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage(raw)}) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + tbl := ppeReadTable(t, data) + assert.Equal(t, "BTC-USDT", ppeCol(t, tbl, "symbol").(*array.String).Value(0)) + assert.Equal(t, int64(1234), ppeCol(t, tbl, "ts").(*array.Int64).Value(0)) +} diff --git a/resources/docker/Dockerfile b/resources/docker/Dockerfile index 23a631ab3b..644b4197ec 100644 --- a/resources/docker/Dockerfile +++ b/resources/docker/Dockerfile @@ -1,7 +1,11 @@ -FROM golang:1.27.1 AS build +FROM --platform=$BUILDPLATFORM golang:1.27.1 AS build +# Cross-compile natively for the target platform instead of emulating it. +ARG TARGETOS=linux +ARG TARGETARCH ENV CGO_ENABLED=0 -ENV GOOS=linux +ENV GOOS=$TARGETOS +ENV GOARCH=$TARGETARCH RUN useradd -u 10001 bento WORKDIR /go/src/github.com/warpstreamlabs/bento/ diff --git a/website/docs/components/inputs/kafka_franz.md b/website/docs/components/inputs/kafka_franz.md index 55d213bfd2..1be209a82b 100644 --- a/website/docs/components/inputs/kafka_franz.md +++ b/website/docs/components/inputs/kafka_franz.md @@ -75,6 +75,7 @@ input: client_certs: [] sasl: [] # No default (optional) multi_header: false + add_record_metadata: true batching: count: 0 byte_size: 0 @@ -770,6 +771,14 @@ Decode headers into lists to allow handling of multiple values with the same key Type: `bool` Default: `false` +### `add_record_metadata` + +Add the `kafka_*` metadata fields and record headers to each message. Disabling this avoids several allocations per record when no downstream component reads them. + + +Type: `bool` +Default: `true` + ### `batching` Allows you to configure a [batching policy](/docs/configuration/batching) that applies to individual topic partitions in order to batch messages together before flushing them for processing. Batching can be beneficial for performance as well as useful for windowed processing, and doing so this way preserves the ordering of topic partitions. diff --git a/website/docs/components/processors/protobuf_parquet_encode.md b/website/docs/components/processors/protobuf_parquet_encode.md new file mode 100644 index 0000000000..988060c8a5 --- /dev/null +++ b/website/docs/components/processors/protobuf_parquet_encode.md @@ -0,0 +1,238 @@ +--- +title: protobuf_parquet_encode +slug: protobuf_parquet_encode +type: processor +status: beta +categories: ["Parsing"] +--- + + + +import Tabs from '@theme/Tabs'; +import TabItem from '@theme/TabItem'; + +:::caution BETA +This component is mostly stable but breaking changes could still be made outside of major version releases if a fundamental problem with the component is found. +::: +Decodes a batch of protobuf messages straight into columnar Parquet files, without an intermediate JSON or structured representation. + + + + + + +```yml +# Common config fields, showing default values +label: "" +protobuf_parquet_encode: + message: "" # No default (required) + import_paths: [] + columns: [] # No default (required) + partition: + field: "" # No default (required) + unit: us + layout: year=2006/month=01/day=02/hour=15 # No default (required) + metadata_key: partition + compression: zstd + skip_invalid_messages: true +``` + + + + +```yml +# All config fields, showing default values +label: "" +protobuf_parquet_encode: + message: "" # No default (required) + import_paths: [] + columns: [] # No default (required) + partition: + field: "" # No default (required) + unit: us + layout: year=2006/month=01/day=02/hour=15 # No default (required) + metadata_key: partition + compression: zstd + compression_level: 3 + dictionary: true + skip_invalid_messages: true +``` + + + + +Each message in the batch must be a single serialized protobuf message of the configured type. The wire format is decoded directly into typed Apache Arrow column builders and written as a Parquet file with one row group, which is an order of magnitude cheaper than chaining the `protobuf` (to_json), `mapping` and `parquet_encode` processors. + +Columns are emitted in the configured order. A column either copies a protobuf field (optionally renamed) or holds a constant string. Proto3 fields that are absent from a message take their zero value. Repeated scalar fields become Parquet LIST columns and string columns are dictionary encoded when `dictionary` is enabled. + +When `partition` is set, the batch is split by formatting a timestamp field with a Go time layout, producing one Parquet file per distinct key, and the key is written to the configured metadata field of each output message. Every output message otherwise carries the metadata of the first input message of its partition. + +Messages that cannot be decoded are dropped and logged when `skip_invalid_messages` is true, and fail the batch otherwise. + +## Examples + + + + + +Batches records at the output, splits each batch by hour of a microsecond timestamp and uploads one Parquet file per hour. + +```yaml +output: + aws_s3: + bucket: raw-archive + path: 'quotes/${! meta("partition") }/${! uuid_v4() }.parquet' + batching: + count: 50000 + period: 60s + processors: + - protobuf_parquet_encode: + message: live.producers.v1.QuoteRecord + import_paths: [ /etc/bento/proto ] + columns: + - { name: exchange, constant: binance-futures } + - { name: symbol } + - { name: time } + - { name: local_timestamp } + - { name: bid_price } + partition: + field: local_timestamp + unit: us + layout: year=2006/month=01/day=02/hour=15 +``` + + + + +## Fields + +### `message` + +The fully qualified name of the protobuf message. + + +Type: `string` + +### `import_paths` + +Directories to load `.proto` files from. + + +Type: `array` +Default: `[]` + +### `columns` + +The output columns, in order. + + +Type: `array` + +### `columns[].name` + +The Parquet column name. + + +Type: `string` + +### `columns[].field` + +The protobuf field copied into the column. Defaults to the column name. + + +Type: `string` + +### `columns[].constant` + +A constant string value for every row, instead of a protobuf field. + + +Type: `string` + +### `partition` + +Split the batch into one Parquet file per formatted timestamp key. + + +Type: `object` + +### `partition.field` + +An integer protobuf field holding a Unix timestamp. + + +Type: `string` + +### `partition.unit` + +The unit of the timestamp field. + + +Type: `string` +Default: `"us"` +Options: `s`, `ms`, `us`, `ns`. + +### `partition.layout` + +A Go time layout, evaluated in UTC, that forms the partition key. + + +Type: `string` + +```yml +# Examples + +layout: year=2006/month=01/day=02/hour=15 +``` + +### `partition.metadata_key` + +The metadata field that receives the partition key. + + +Type: `string` +Default: `"partition"` + +### `compression` + +Sorry! This field is missing documentation. + + +Type: `string` +Default: `"zstd"` +Options: `zstd`, `snappy`, `gzip`, `lz4raw`, `uncompressed`. + +### `compression_level` + +The compression level, for codecs that support one. + + +Type: `int` +Default: `3` + +### `dictionary` + +Dictionary encode string columns. + + +Type: `bool` +Default: `true` + +### `skip_invalid_messages` + +Drop (and log) messages that cannot be decoded instead of failing the batch. + + +Type: `bool` +Default: `true` + +