diff --git a/.github/workflows/aperiodic_test.yml b/.github/workflows/aperiodic_test.yml new file mode 100644 index 0000000000..a6664c43d3 --- /dev/null +++ b/.github/workflows/aperiodic_test.yml @@ -0,0 +1,116 @@ +name: Aperiodic Test + +# Tests and lints only the Go packages a pull request touches. The upstream Test +# workflow runs the whole module and is disabled on this fork. + +on: + pull_request: + workflow_dispatch: + inputs: + packages: + description: 'Space-separated package directories to test (e.g. "./internal/impl/parquet")' + required: true + +concurrency: + group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + +permissions: + contents: read + +jobs: + changes: + if: github.repository == 'aperiodic-io/bento' + runs-on: ubicloud-standard-2 + outputs: + packages: ${{ steps.packages.outputs.packages }} + modules: ${{ steps.packages.outputs.modules }} + steps: + - uses: actions/checkout@v7 + with: + # The pull request's merge commit and its first parent, the base. + fetch-depth: 2 + + - name: Find touched packages + id: packages + env: + INPUT_PACKAGES: ${{ inputs.packages }} + run: | + if [ "$GITHUB_EVENT_NAME" = workflow_dispatch ]; then + echo "packages=$INPUT_PACKAGES" >> "$GITHUB_OUTPUT" + exit 0 + fi + changed="$(git diff --name-only HEAD^1 HEAD)" + echo "$changed" + # Each changed file belongs to the package of the nearest directory + # holding Go files; testdata belongs to the package that owns it. + packages="$(echo "$changed" | while read -r f; do + d="$(dirname "$f")" + d="${d%%/testdata*}" + while [ "$d" != . ] && ! compgen -G "$d/*.go" > /dev/null; do + d="$(dirname "$d")" + done + if [ "$d" != . ]; then echo "./$d"; fi + done | sort -u | tr '\n' ' ')" + echo "packages=$packages" >> "$GITHUB_OUTPUT" + if grep -qxE 'go\.(mod|sum)' <<< "$changed"; then + echo "modules=true" >> "$GITHUB_OUTPUT" + fi + echo "Touched packages: ${packages:-none}" + + test: + needs: changes + if: needs.changes.outputs.packages != '' || needs.changes.outputs.modules == 'true' + runs-on: ubicloud-standard-2 + env: + PACKAGES: ${{ needs.changes.outputs.packages }} + steps: + - uses: actions/checkout@v7 + + - uses: actions/setup-go@v7 + with: + go-version-file: go.mod + check-latest: true + + - name: Deps + if: needs.changes.outputs.modules == 'true' + run: go mod tidy && git diff --exit-code go.mod go.sum || { >&2 echo "Stale go.{mod,sum} detected. This can be fixed with 'make deps'."; exit 1; } + + # Downloads only the modules the touched packages build with, retrying the + # module proxy's intermittent stream errors. + - name: Download modules + if: env.PACKAGES != '' + run: | + for attempt in 1 2 3 4 5; do + go list -deps -test $PACKAGES > /dev/null && exit 0 + sleep $((attempt * 10)) + done + exit 1 + + - name: Vet + if: env.PACKAGES != '' + run: go vet $PACKAGES + + - name: Test + if: env.PACKAGES != '' + run: go test -timeout 10m $PACKAGES + + golangci-lint: + needs: changes + if: needs.changes.outputs.packages != '' + runs-on: ubicloud-standard-2 + env: + CGO_ENABLED: 0 + steps: + - uses: actions/checkout@v7 + + - uses: actions/setup-go@v7 + with: + go-version-file: go.mod + check-latest: true + + - name: Lint + uses: golangci/golangci-lint-action@v9 + with: + version: v2.13.2 + args: ${{ needs.changes.outputs.packages }} diff --git a/internal/impl/kafka/input_kafka_franz.go b/internal/impl/kafka/input_kafka_franz.go index 2da0ceca9f..929bc4531e 100644 --- a/internal/impl/kafka/input_kafka_franz.go +++ b/internal/impl/kafka/input_kafka_franz.go @@ -1,6 +1,7 @@ package kafka import ( + "bytes" "context" "crypto/tls" "errors" @@ -167,6 +168,7 @@ With this option, you can return topic order and per-topic partition ordering. T 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.NewBoolField("copy_record_values").Description("Copy each record's value into its message instead of referencing the fetch response it arrived in. A record's key, value and header values are slices of that response, which holds every record fetched from the broker in the same request, so a message held for long (e.g. in a large output batch) otherwise keeps the whole response in memory. Costs one allocation per record.").Default(false).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()). @@ -229,6 +231,7 @@ type franzKafkaReader struct { regexPattern bool multiHeader bool addMetadata bool + copyValues bool batchPolicy service.BatchPolicy reconnectOnUnknownTopic bool @@ -505,6 +508,9 @@ func newFranzKafkaReaderFromConfig(conf *service.ParsedConfig, res *service.Reso if f.addMetadata, err = conf.FieldBool("add_record_metadata"); err != nil { return nil, err } + if f.copyValues, err = conf.FieldBool("copy_record_values"); err != nil { + return nil, err + } if f.multiHeader, err = conf.FieldBool("multi_header"); err != nil { return nil, err } @@ -521,10 +527,17 @@ type msgWithRecord struct { } func (f *franzKafkaReader) recordToMessage(record *kgo.Record) *msgWithRecord { - msg := service.NewMessage(record.Value) + value := record.Value + if f.copyValues && value != nil { + value = bytes.Clone(value) + } + msg := service.NewMessage(value) if !f.addMetadata { record.Key = nil record.Value = nil + if f.copyValues { + record.Headers = nil + } return &msgWithRecord{msg: msg, r: record} } if record.Key != nil { @@ -558,6 +571,11 @@ func (f *franzKafkaReader) recordToMessage(record *kgo.Record) *msgWithRecord { // potentially be a source of problems so treat this as sus. record.Key = nil record.Value = nil + if f.copyValues { + // header values alias the fetch response too; they were copied into + // metadata above + record.Headers = nil + } return &msgWithRecord{ msg: msg, diff --git a/internal/impl/kafka/input_kafka_franz_metadata_test.go b/internal/impl/kafka/input_kafka_franz_metadata_test.go index d7e7eb0841..c0a9971f1a 100644 --- a/internal/impl/kafka/input_kafka_franz_metadata_test.go +++ b/internal/impl/kafka/input_kafka_franz_metadata_test.go @@ -39,3 +39,43 @@ func TestFranzRecordToMessageMetadata(t *testing.T) { require.NoError(t, bare.msg.MetaWalkMut(func(string, any) error { count++; return nil })) assert.Zero(t, count) } + +func TestFranzRecordToMessageCopyRecordValues(t *testing.T) { + // 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 in a long batch (and the + // record kept for checkpointing) would keep that whole response alive: with + // copy_record_values neither may reference it. + for _, addMetadata := range []bool{true, false} { + fetched := []byte("payload|k|abc") + record := &kgo.Record{ + Value: fetched[0:7], Key: fetched[8:9], Topic: "metric.v1.x.15s", Partition: 1, Offset: 42, + Headers: []kgo.RecordHeader{{Key: "h", Value: fetched[10:13]}}, + } + m := (&franzKafkaReader{addMetadata: addMetadata, copyValues: true}).recordToMessage(record) + for i := range fetched { + fetched[i] = 'X' // the fetch buffer is reused or freed once nothing points at it + } + b, err := m.msg.AsBytes() + require.NoError(t, err) + assert.Equal(t, "payload", string(b), "addMetadata=%v: the message must own its bytes", addMetadata) + assert.Nil(t, m.r.Key, "addMetadata=%v", addMetadata) + assert.Nil(t, m.r.Value, "addMetadata=%v", addMetadata) + assert.Nil(t, m.r.Headers, "addMetadata=%v: header values alias the fetch buffer", addMetadata) + assert.Equal(t, int64(42), m.r.Offset, "addMetadata=%v: checkpointing needs the offset", addMetadata) + if addMetadata { + got, _ := m.msg.MetaGetMut("h") + assert.Equal(t, "abc", got) + got, _ = m.msg.MetaGetMut("kafka_key") + assert.Equal(t, "k", got) + } + } + + // The default is unchanged: the message shares the record's bytes. + fetched := []byte("payload") + m := (&franzKafkaReader{addMetadata: true}).recordToMessage(&kgo.Record{Value: fetched}) + fetched[0] = 'X' + b, err := m.msg.AsBytes() + require.NoError(t, err) + assert.Equal(t, "Xayload", string(b)) +} diff --git a/internal/impl/parquet/processor_json_encode.go b/internal/impl/parquet/processor_json_encode.go new file mode 100644 index 0000000000..d7dd7ef964 --- /dev/null +++ b/internal/impl/parquet/processor_json_encode.go @@ -0,0 +1,632 @@ +package parquet + +import ( + "bytes" + "context" + "encoding/json" + "encoding/json/jsontext" + "errors" + "fmt" + "io" + "math" + "reflect" + "strings" + + "github.com/parquet-go/parquet-go" + + "github.com/warpstreamlabs/bento/internal/value" + "github.com/warpstreamlabs/bento/public/service" +) + +const ( + jpeFieldColumns = "columns" + jpeFieldColumnName = "name" + jpeFieldColumnValue = "value" + jpeFieldNaNForNull = "nan_for_null" + jpeFieldPartition = "partition" + jpeFieldPartitionPath = "path" + jpeFieldPartitionUnit = "time_unit" + jpeFieldPartitionMetaKey = "metadata_key" +) + +func jsonParquetEncodeSpec() *service.ConfigSpec { + return service.NewConfigSpec(). + Beta(). + Categories("Parsing"). + Summary("Encodes a batch of flat JSON objects straight into Parquet, row by row, without holding the batch as structured messages."). + Description(` +Produces the same Parquet as chaining a `+"`mapping`"+` that coerces each column, a `+"`catch`"+` that drops the rows it rejects, `+"`group_by_value`"+` on a partition path and `+"`parquet_encode`"+`, at a fraction of the memory: a batch is held only as its raw bytes until it is processed, each message is then parsed once into a Parquet row, and the files are encoded one at a time. + +The schema is `+"`parquet_encode`"+`'s, restricted to flat `+"`UTF8`"+`, `+"`INT32`"+`, `+"`INT64`"+`, `+"`FLOAT`"+` and `+"`DOUBLE`"+` columns, and the file is written by the same encoder. A column reads the message field of its name and coerces it as Bloblang would: + +- `+"`UTF8`"+`: `+"`.not_null().string()`"+`, so a number keeps its literal text; +- `+"`INT32`"+`, `+"`INT64`"+`: `+"`.int32()`"+`, `+"`.int64()`"+` (a fraction, an exponent or an out-of-range value is rejected; a string is parsed); +- `+"`FLOAT`"+`, `+"`DOUBLE`"+`: `+"`.number()`"+`, narrowed to 32 bits for `+"`FLOAT`"+`. + +A required column the message lacks rejects it. A `+"`null`"+` in an optional column is written as NULL, and in a required float column as NaN when `+"`nan_for_null`"+` is set; otherwise it rejects the message. A column may instead take an interpolated `+"`value`"+`, evaluated per message against its metadata only. + +A message that is not a single JSON object, or that any column rejects, is dropped: it is logged, counted in `+"`json_parquet_encode_dropped`"+` and acknowledged with the batch. + +When `+"`partition`"+` is set the batch is split into one Parquet file per distinct partition path, in order of first appearance, and the path is written to the configured metadata key. `+"`{column}`"+` in the path is replaced by the column's value as Bloblang's `+"`format`"+` would print the field, and `+"`{column|layout}`"+` by the column's number, read in `+"`time_unit`"+`, as a UTC time in the Go layout. Every output message carries the metadata of the first message of its partition.`). + Field(parquetSchemaConfig()). + Field(service.NewObjectListField(jpeFieldColumns, + service.NewStringField(jpeFieldColumnName).Description("The schema column this sets."), + service.NewInterpolatedStringField(jpeFieldColumnValue).Description("The column's value, from the message's metadata. It must not reference the message content."), + ).Description("Columns set from an interpolation instead of the message field of their name.").Default([]any{})). + Field(service.NewBoolField(jpeFieldNaNForNull).Description("Write a `null` in a required `FLOAT` or `DOUBLE` column as NaN instead of dropping the message.").Default(false)). + Field(service.NewObjectField(jpeFieldPartition, + service.NewStringField(jpeFieldPartitionPath).Description("The partition path, with `{column}` and `{column|layout}` placeholders.").Example("day={time|2006-01-02}/{symbol}"), + service.NewStringEnumField(jpeFieldPartitionUnit, "s", "ms", "us", "ns").Description("The unit of numbers formatted as a time.").Default("us"), + service.NewStringField(jpeFieldPartitionMetaKey).Description("The metadata field that receives the partition path.").Default("partition"), + ).Description("Split the batch into one Parquet file per partition path.").Optional()). + Field(service.NewStringEnumField("default_compression", + "uncompressed", "snappy", "gzip", "brotli", "zstd", "lz4raw", + ).Description("The default compression type to use for fields.").Default("uncompressed")). + Field(service.NewStringEnumField("default_encoding", + "DELTA_LENGTH_BYTE_ARRAY", "PLAIN", "RLE_DICTIONARY", + ).Description("The default encoding type to use for fields.").Default("DELTA_LENGTH_BYTE_ARRAY").Advanced()). + Example("Archiving JSON rows to S3 as daily Parquet", + "Batches rows at the output and uploads one Parquet file per exchange and UTC day of a microsecond timestamp.", + ` +output: + aws_s3: + bucket: metrics + path: 'ohlcv/${! meta("partition") }/${! uuid_v4() }.parquet' + batching: + byte_size: 10485760 + period: 30s + processors: + - json_parquet_encode: + default_compression: zstd + nan_for_null: true + schema: + - { name: exchange, type: UTF8 } + - { name: symbol, type: UTF8 } + - { name: time, type: INT64 } + - { name: close, type: DOUBLE } + - { name: volume, type: DOUBLE, optional: true } + columns: + - { name: exchange, value: '${! @kafka_topic.split(".").index(2) }' } + partition: + path: 'exchange={exchange}/{time|year=2006/month=01/day=02}' + time_unit: us +`) +} + +func init() { + err := service.RegisterBatchProcessor("json_parquet_encode", jsonParquetEncodeSpec(), + func(conf *service.ParsedConfig, mgr *service.Resources) (service.BatchProcessor, error) { + return newJSONParquetEncoder(conf, mgr) + }) + if err != nil { + panic(err) + } +} + +//------------------------------------------------------------------------------ + +type jpeKind int + +const ( + jpeUTF8 jpeKind = iota + jpeInt32 + jpeInt64 + jpeFloat + jpeDouble +) + +type jpeColumn struct { + name string + kind jpeKind + optional bool + leaf int // index of the Parquet leaf column + interp *service.InterpolatedString // set when the column comes from metadata +} + +// jpePart is one piece of a partition path: literal text, or a column formatted +// as Bloblang's format prints it (layout empty) or as a time (layout set). +type jpePart struct { + literal string + column int + layout string +} + +type jsonParquetEncoder struct { + log *service.Logger + mDropped *service.MetricCounter + + schema *parquet.Schema + codec parquet.WriterOption + columns []jpeColumn + byName map[string]int + nanNull bool + + partition []jpePart + divisor float64 + partitionMeta string +} + +func newJSONParquetEncoder(conf *service.ParsedConfig, mgr *service.Resources) (*jsonParquetEncoder, error) { + e := &jsonParquetEncoder{ + log: mgr.Logger(), + mDropped: mgr.Metrics().NewCounter("json_parquet_encode_dropped"), + byName: map[string]int{}, + } + + fields, err := conf.FieldObjectList("schema") + if err != nil { + return nil, err + } + if len(fields) == 0 { + return nil, errors.New("at least one schema column is required") + } + for _, f := range fields { + name, err := f.FieldString("name") + if err != nil { + return nil, err + } + if children, err := f.FieldAnyList("fields"); err == nil && len(children) > 0 { + return nil, fmt.Errorf("column '%v': only flat columns are supported", name) + } + typ, err := f.FieldString("type") + if err != nil { + return nil, fmt.Errorf("column '%v': %w", name, err) + } + kind, ok := map[string]jpeKind{"UTF8": jpeUTF8, "INT32": jpeInt32, "INT64": jpeInt64, "FLOAT": jpeFloat, "DOUBLE": jpeDouble}[typ] + if !ok { + return nil, fmt.Errorf("column '%v': type %v is not supported, only UTF8, INT32, INT64, FLOAT and DOUBLE", name, typ) + } + if repeated, _ := f.FieldBool("repeated"); repeated { + return nil, fmt.Errorf("column '%v': repeated columns are not supported", name) + } + optional, err := f.FieldBool("optional") + if err != nil { + return nil, err + } + if _, exists := e.byName[name]; exists { + return nil, fmt.Errorf("column '%v' is defined twice", name) + } + e.byName[name] = len(e.columns) + e.columns = append(e.columns, jpeColumn{name: name, kind: kind, optional: optional}) + } + + overrides, err := conf.FieldObjectList(jpeFieldColumns) + if err != nil { + return nil, err + } + for _, o := range overrides { + name, err := o.FieldString(jpeFieldColumnName) + if err != nil { + return nil, err + } + i, ok := e.byName[name] + if !ok { + return nil, fmt.Errorf("columns: '%v' is not a schema column", name) + } + if e.columns[i].kind != jpeUTF8 { + return nil, fmt.Errorf("columns: '%v' must be a UTF8 column to take an interpolated value", name) + } + if e.columns[i].interp, err = o.FieldInterpolatedString(jpeFieldColumnValue); err != nil { + return nil, err + } + } + if e.nanNull, err = conf.FieldBool(jpeFieldNaNForNull); err != nil { + return nil, err + } + + // The schema and writer options are parquet_encode's, built by the same code, + // so the files cannot drift from what it writes. + encodingStr, err := conf.FieldString("default_encoding") + if err != nil { + return nil, err + } + encodingTag, ok := map[string]string{ + parquet.Plain.String(): "plain", parquet.DeltaLengthByteArray.String(): "delta", parquet.RLEDictionary.String(): "dict", + }[encodingStr] + if !ok { + return nil, fmt.Errorf("default_encoding type %v not recognised", encodingStr) + } + schemaType, err := GenerateStructType(conf, schemaOpts{optionalsAsStructTags: true, defaultEncoding: encodingTag}) + if err != nil { + return nil, fmt.Errorf("failed to generate struct type from parquet schema: %w", err) + } + e.schema = parquet.SchemaOf(reflect.New(schemaType).Interface()) + for leaf, path := range e.schema.Columns() { + if len(path) != 1 { + return nil, fmt.Errorf("unexpected nested column %v", path) + } + i, ok := e.byName[path[0]] + if !ok { + return nil, fmt.Errorf("schema column %v has no definition", path[0]) + } + e.columns[i].leaf = leaf + } + + compressStr, err := conf.FieldString("default_compression") + if err != nil { + return nil, err + } + codec, ok := map[string]parquet.WriterOption{ + "uncompressed": parquet.Compression(&parquet.Uncompressed), "snappy": parquet.Compression(&parquet.Snappy), + "gzip": parquet.Compression(&parquet.Gzip), "brotli": parquet.Compression(&parquet.Brotli), + "zstd": parquet.Compression(&parquet.Zstd), "lz4raw": parquet.Compression(&parquet.Lz4Raw), + }[compressStr] + if !ok { + return nil, fmt.Errorf("default_compression type %v not recognised", compressStr) + } + e.codec = codec + + if conf.Contains(jpeFieldPartition) { + pc := conf.Namespace(jpeFieldPartition) + path, err := pc.FieldString(jpeFieldPartitionPath) + if err != nil { + return nil, err + } + if e.partition, err = e.parsePartitionPath(path); err != nil { + return nil, fmt.Errorf("partition: %w", err) + } + unit, err := pc.FieldString(jpeFieldPartitionUnit) + if err != nil { + return nil, err + } + divisor, ok := map[string]float64{"s": 1, "ms": 1e3, "us": 1e6, "ns": 1e9}[unit] + if !ok { + return nil, fmt.Errorf("partition: unknown time_unit '%v'", unit) + } + e.divisor = divisor + if e.partitionMeta, err = pc.FieldString(jpeFieldPartitionMetaKey); err != nil { + return nil, err + } + } + return e, nil +} + +func (e *jsonParquetEncoder) parsePartitionPath(path string) ([]jpePart, error) { + var parts []jpePart + for len(path) > 0 { + open := strings.IndexByte(path, '{') + if open < 0 { + parts = append(parts, jpePart{literal: path, column: -1}) + break + } + if open > 0 { + parts = append(parts, jpePart{literal: path[:open], column: -1}) + } + end := strings.IndexByte(path[open:], '}') + if end < 0 { + return nil, fmt.Errorf("unclosed '{' in %q", path) + } + name, layout, isTime := strings.Cut(path[open+1:open+end], "|") + i, ok := e.byName[name] + if !ok { + return nil, fmt.Errorf("placeholder {%v} is not a schema column", name) + } + if isTime && layout == "" { + return nil, fmt.Errorf("placeholder {%v|}: empty time layout", name) + } + if isTime && e.columns[i].interp != nil { + return nil, fmt.Errorf("placeholder {%v|%v}: an interpolated column cannot be formatted as a time", name, layout) + } + parts = append(parts, jpePart{column: i, layout: layout}) + path = path[open+end+1:] + } + return parts, nil +} + +//------------------------------------------------------------------------------ + +// jpeRow holds the fields of one message the schema reads, as the Go values +// Bento's own JSON decoding gives them (string, json.Number, bool, nil, or a +// map or slice for a nested value), so coercion runs the very functions +// Bloblang's methods do. +type jpeRow struct { + present []bool + values []any +} + +// jpeDecodeOptions make jsontext read a document as encoding/json does, and so +// as Bento's decoding of every message: the last of a duplicated name wins and +// invalid UTF-8 is mangled into U+FFFD. +var jpeDecodeOptions = []jsontext.Options{ + jsontext.AllowDuplicateNames(true), + jsontext.AllowInvalidUTF8(true), +} + +// parse reads one message into r. It fails where Bento's own decoding of the +// message would: anything but exactly one JSON document. A document that is not +// an object fails too, which only drops what the coercion would have dropped +// anyway, since a schema read from the message has at least one column. +func (e *jsonParquetEncoder) parse(dec *jsontext.Decoder, src *bytes.Reader, raw []byte, r *jpeRow) error { + clear(r.present) + clear(r.values) + src.Reset(raw) + dec.Reset(src, jpeDecodeOptions...) + + tok, err := dec.ReadToken() + if err != nil { + return fmt.Errorf("not JSON: %w", err) + } + if tok.Kind() != '{' { + return fmt.Errorf("not a JSON object but %v", tok.Kind()) + } + for dec.PeekKind() != '}' { + name, err := dec.ReadToken() + if err != nil { + return fmt.Errorf("not JSON: %w", err) + } + i, ok := e.byName[name.String()] + if !ok || e.columns[i].interp != nil { + if err := dec.SkipValue(); err != nil { + return fmt.Errorf("not JSON: %w", err) + } + continue + } + if r.values[i], err = jpeGoValue(dec); err != nil { + return fmt.Errorf("not JSON: %w", err) + } + r.present[i] = true + } + if _, err := dec.ReadToken(); err != nil { // the closing '}' + return fmt.Errorf("not JSON: %w", err) + } + if _, err := dec.ReadToken(); !errors.Is(err, io.EOF) { + if err == nil { + return errors.New("message contains multiple valid documents") + } + return fmt.Errorf("not JSON: %w", err) + } + return nil +} + +// jpeGoValue reads the next JSON value as Bento's decoding gives it. +func jpeGoValue(dec *jsontext.Decoder) (any, error) { + switch dec.PeekKind() { + case '{', '[': + v, err := dec.ReadValue() + if err != nil { + return nil, err + } + nested := json.NewDecoder(bytes.NewReader(v)) + nested.UseNumber() + var out any + err = nested.Decode(&out) + return out, err + } + tok, err := dec.ReadToken() + if err != nil { + return nil, err + } + switch tok.Kind() { + case '"': + return tok.String(), nil + case '0': + return json.Number(tok.String()), nil + case 't': + return true, nil + case 'f': + return false, nil + } + return nil, nil +} + +// value coerces column i of r, as the column's Bloblang coercion would, into +// the Parquet value of its leaf. +func (e *jsonParquetEncoder) value(c *jpeColumn, present bool, v any) (parquet.Value, error) { + def := 0 + if c.optional { + def = 1 + } + if !present || v == nil { + switch { + case c.optional: + return parquet.NullValue().Level(0, 0, c.leaf), nil + case !present: + return parquet.Value{}, errors.New("missing") + case e.nanNull && (c.kind == jpeFloat || c.kind == jpeDouble): + v = math.NaN() + default: + return parquet.Value{}, errors.New("value is null") + } + } + var pv parquet.Value + switch c.kind { + case jpeUTF8: + pv = parquet.ByteArrayValue([]byte(jpeString(v))) + case jpeInt32: + i, err := value.IToInt32(v) + if err != nil { + return parquet.Value{}, err + } + pv = parquet.Int32Value(i) + case jpeInt64: + i, err := value.IToInt(v) + if err != nil { + return parquet.Value{}, err + } + pv = parquet.Int64Value(i) + case jpeFloat: + f, err := value.IToNumber(v) + if err != nil { + return parquet.Value{}, err + } + pv = parquet.FloatValue(float32(f)) + case jpeDouble: + f, err := value.IToNumber(v) + if err != nil { + return parquet.Value{}, err + } + pv = parquet.DoubleValue(f) + } + return pv.Level(0, def, c.leaf), nil +} + +// jpeString is Bloblang's .string(), whose objects are marshalled by gabs. +func jpeString(v any) string { + return value.IToString(v) +} + +// partitionPath formats the partition path of a row, as the Bloblang +// "...".format(this.a, ...) and (this.t / divisor).ts_format(layout, "UTC") +// it replaces would. +func (e *jsonParquetEncoder) partitionPath(b *strings.Builder, r *jpeRow, interpolated []string) error { + b.Reset() + for _, p := range e.partition { + switch { + case p.column < 0: + b.WriteString(p.literal) + case e.columns[p.column].interp != nil: + b.WriteString(interpolated[p.column]) + case p.layout == "": + fmt.Fprintf(b, "%s", r.values[p.column]) + default: + n, err := value.IGetNumber(r.values[p.column]) + if err != nil { + return fmt.Errorf("partition: %v: %w", e.columns[p.column].name, err) + } + t, err := value.IGetTimestamp(n / e.divisor) + if err != nil { + return fmt.Errorf("partition: %v: %w", e.columns[p.column].name, err) + } + b.WriteString(t.UTC().Format(p.layout)) + } + } + return nil +} + +//------------------------------------------------------------------------------ + +// jpeGroup is one output file: the rows of one partition path, held until the +// batch is read, and then encoded while no other group's file is. +type jpeGroup struct { + first *service.Message + key string + rows []parquet.Row +} + +func (e *jsonParquetEncoder) ProcessBatch(ctx context.Context, batch service.MessageBatch) ([]service.MessageBatch, error) { + if len(batch) == 0 { + return nil, nil + } + var ( + groups []*jpeGroup + byKey = map[string]*jpeGroup{} + current *jpeGroup + src bytes.Reader + dec = jsontext.NewDecoder(&src, jpeDecodeOptions...) + row = &jpeRow{present: make([]bool, len(e.columns)), values: make([]any, len(e.columns))} + interpolated = make([]string, len(e.columns)) + key strings.Builder + dropped int + ) + drop := func(err error) { + dropped++ + e.log.Errorf("json_parquet_encode: dropping a message that is not a valid row: %v", err) + } + for _, msg := range batch { + raw, err := msg.AsBytes() + if err == nil && len(raw) == 0 { + err = errors.New("empty message") + } + if err == nil { + err = e.parse(dec, &src, raw, row) + } + if err == nil { + err = e.interpolate(msg, interpolated) + } + if err != nil { + drop(err) + continue + } + values := make(parquet.Row, len(e.columns)) + for i := range e.columns { + c := &e.columns[i] + if c.interp != nil { + values[c.leaf] = parquet.ByteArrayValue([]byte(interpolated[i])).Level(0, 0, c.leaf) + continue + } + if values[c.leaf], err = e.value(c, row.present[i], row.values[i]); err != nil { + err = fmt.Errorf("%v: %w", c.name, err) + break + } + } + if err != nil { + drop(err) + continue + } + k := "" + if e.partition != nil { + if err := e.partitionPath(&key, row, interpolated); err != nil { + drop(err) + continue + } + k = key.String() + } + if current == nil || current.key != k { + if current = byKey[k]; current == nil { + current = &jpeGroup{first: msg, key: k} + byKey[k] = current + groups = append(groups, current) + } + } + current.rows = append(current.rows, values) + } + if dropped > 0 { + e.mDropped.Incr(int64(dropped)) + } + if len(groups) == 0 { + return nil, nil + } + + // Files are encoded one at a time, as parquet_encode encodes each group + // group_by_value hands it: a writer's column buffers are the bulk of what + // encoding holds, so a batch of many partitions must not hold one per file. + // Each file takes a new writer, since a reset one does not write the same + // bytes as a new one. + out := make(service.MessageBatch, 0, len(groups)) + for _, g := range groups { + var buf bytes.Buffer + if err := jpeEncode(parquet.NewGenericWriter[any](&buf, e.schema, e.codec), g.rows); err != nil { + return nil, err + } + g.rows = nil + m := g.first.Copy() + m.SetBytes(buf.Bytes()) + if e.partition != nil { + m.MetaSetMut(e.partitionMeta, g.key) + } + out = append(out, m) + } + return []service.MessageBatch{out}, nil +} + +// interpolate resolves the interpolated columns of a message. +func (e *jsonParquetEncoder) interpolate(msg *service.Message, interpolated []string) error { + for i, c := range e.columns { + if c.interp == nil { + continue + } + s, err := c.interp.TryString(msg) + if err != nil { + return fmt.Errorf("%v: %w", c.name, err) + } + interpolated[i] = s + } + return nil +} + +// jpeEncode writes rows as one whole file. +func jpeEncode(w *parquet.GenericWriter[any], rows []parquet.Row) (err error) { + defer func() { + if r := recover(); r != nil { + err = fmt.Errorf("encoding panic: %v", r) + } + }() + if _, err = w.WriteRows(rows); err != nil { + return err + } + return w.Close() +} + +func (e *jsonParquetEncoder) Close(ctx context.Context) error { + return nil +} diff --git a/internal/impl/parquet/processor_json_encode_bench_test.go b/internal/impl/parquet/processor_json_encode_bench_test.go new file mode 100644 index 0000000000..459d887975 --- /dev/null +++ b/internal/impl/parquet/processor_json_encode_bench_test.go @@ -0,0 +1,131 @@ +package parquet_test + +import ( + "bytes" + "math/rand/v2" + "runtime" + "runtime/metrics" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/require" + + "github.com/warpstreamlabs/bento/public/bloblang" + "github.com/warpstreamlabs/bento/public/service" +) + +// benchBatch is a batch of valid rows the size an archiver closes a batch at. +func benchBatch(tb testing.TB) []parityInput { + tb.Helper() + r := rand.New(rand.NewPCG(7, 8)) + var in []parityInput + for size := 0; size < 10<<20; { + m := randomRow(r) + if bytes.Contains(m.body, []byte(`"`+missing)) { + continue + } + size += len(m.body) + in = append(in, m) + } + return in +} + +func heapBytes() uint64 { + s := []metrics.Sample{{Name: "/memory/classes/heap/objects:bytes"}} + metrics.Read(s) + return s[0].Value.Uint64() +} + +func liveHeap() uint64 { + runtime.GC() + runtime.GC() + return heapBytes() +} + +// BenchmarkJSONParquetHeldBatch measures what holding a batch costs until it +// closes: the replaced pipeline holds each row as the structured value its +// mapping made, json_parquet_encode the row's bytes. +func BenchmarkJSONParquetHeldBatch(b *testing.B) { + in := benchBatch(b) + raw := 0 + for _, m := range in { + raw += len(m.body) + } + exe, err := bloblang.Parse(archiveSchema.legacyMapping()) + require.NoError(b, err) + + b.Run("legacy", func(b *testing.B) { + for i := 0; i < b.N; i++ { + before := liveHeap() + held := make(service.MessageBatch, 0, len(in)) + for _, m := range in { + msg := service.NewMessage(bytes.Clone(m.body)) + msg.MetaSetMut("kafka_topic", m.topic) + out, err := msg.BloblangQuery(exe) + if err == nil && out != nil { + held = append(held, out) + } + } + after := liveHeap() + b.ReportMetric(float64(after-before)/float64(raw), "held-bytes/JSON-byte") + runtime.KeepAlive(held) + } + }) + b.Run("json_parquet_encode", func(b *testing.B) { + for i := 0; i < b.N; i++ { + before := liveHeap() + held := make(service.MessageBatch, 0, len(in)) + for _, m := range in { + msg := service.NewMessage(bytes.Clone(m.body)) + msg.MetaSetMut("kafka_topic", m.topic) + held = append(held, msg) + } + after := liveHeap() + b.ReportMetric(float64(after-before)/float64(raw), "held-bytes/JSON-byte") + runtime.KeepAlive(held) + } + }) +} + +// BenchmarkJSONParquetEncode runs a full batch through each pipeline, and +// reports the heap's peak above what it held before the batch. +func BenchmarkJSONParquetEncode(b *testing.B) { + in := benchBatch(b) + raw := 0 + for _, m := range in { + raw += len(m.body) + } + p := pairFor(b, archiveSchema) + for name, s := range map[string]*parityStream{"legacy": p.legacy, "json_parquet_encode": p.next} { + b.Run(name, func(b *testing.B) { + b.SetBytes(int64(raw)) + b.ReportAllocs() + var peakSum float64 + for i := 0; i < b.N; i++ { + before := liveHeap() + var peak atomic.Uint64 + done := make(chan struct{}) + go func() { + t := time.NewTicker(time.Millisecond) + defer t.Stop() + for { + select { + case <-done: + return + case <-t.C: + if h := heapBytes(); h > peak.Load() { + peak.Store(h) + } + } + } + }() + files := s.run(b, in) + close(done) + require.NotEmpty(b, files) + peakSum += float64(peak.Load()-min(before, peak.Load())) / (1 << 20) + } + b.ReportMetric(peakSum/float64(b.N), "peak-heap-MiB") + }) + } +} diff --git a/internal/impl/parquet/processor_json_encode_parity_cases_test.go b/internal/impl/parquet/processor_json_encode_parity_cases_test.go new file mode 100644 index 0000000000..2cd7dbc7ab --- /dev/null +++ b/internal/impl/parquet/processor_json_encode_parity_cases_test.go @@ -0,0 +1,270 @@ +package parquet_test + +import ( + "fmt" + "math/rand/v2" + "strings" + "testing" + + "github.com/stretchr/testify/require" +) + +const parityTopic = "metric.v1.binance-futures.15s" + +// parityBase is a valid row of archiveSchema, as JSON fields in order. +var parityBase = [][2]string{ + {"exchange", `11`}, // the row's own exchange id, which the topic's name overrides + {"symbol", `"perpetual-BTC-USDT:USD"`}, + {"interval", `"15s"`}, + {"timestamp_type", `"true"`}, + {"day", `"2026-09-29"`}, + {"time", `1790677155000000`}, + {"count", `10`}, + {"count_opt", `6`}, + {"ratio", `0.16286439`}, + {"ratio_opt", `1.6666666`}, + {"price", `0.012981182352941168`}, + {"price_opt", `3.1113223111729823e-10`}, +} + +// rowWith renders parityBase with field set to raw JSON (missing removes it). +func rowWith(field, raw string) []byte { + var parts []string + for _, kv := range parityBase { + switch { + case kv[0] != field: + parts = append(parts, fmt.Sprintf("%q:%s", kv[0], kv[1])) + case raw != missing: + parts = append(parts, fmt.Sprintf("%q:%s", kv[0], raw)) + } + } + return []byte("{" + strings.Join(parts, ",") + "}") +} + +const missing = "" + +// valueVariants are the JSON values every column is tried with: every kind, +// every numeric literal form and range edge, and the strings Bloblang parses. +var valueVariants = []string{ + missing, `null`, `true`, `false`, + `"abc"`, `""`, `"我踏马来了"`, `"é😀"`, `"a\"b\\c\n\t\u0000"`, `"\ud800"`, "\"\xff\xfe\"", `"<&>"`, + `0`, `-0`, `1`, `-1`, `1.0`, `-0.0`, `1.5`, `-1.5e-3`, `1e2`, `1E+2`, `0.1`, `0.30000000000000004`, + `2147483647`, `2147483648`, `-2147483648`, `-2147483649`, + `9223372036854775807`, `9223372036854775808`, `-9223372036854775808`, `-9223372036854775809`, `12345678901234567890`, + `1e308`, `1.7976931348623157e308`, `1e309`, `1e400`, `-1e400`, `5e-324`, `1e-400`, + `3.4028234663852886e38`, `3.4028236e38`, `1.401298464324817e-45`, `1e-46`, + `1790677155000000`, `1790677155999999`, `-1`, + `"123"`, `"-5"`, `"0x10"`, `"010"`, `"0b101"`, `"1_000"`, `" 1"`, `"1 "`, `"1.5"`, `"1e2"`, `"NaN"`, `"nan"`, `"inf"`, `"-Infinity"`, `"+1"`, `"9223372036854775808"`, + `{}`, `{"a":1}`, `{"z":"<&>","b":1.50,"a":[true,null]}`, `[]`, `[1,"x",{"k":null}]`, +} + +// rowsIn counts the rows of the files, so a case cannot pass with both sides +// writing nothing. +func rowsIn(t testing.TB, files []parityFile) int { + n := 0 + for _, f := range files { + n += len(parityRows(t, f.data)) + } + return n +} + +func TestJSONParquetParityEveryColumnEveryValue(t *testing.T) { + base := requireParity(t, archiveSchema, []parityInput{{topic: parityTopic, body: rowWith("", "")}}) + require.Equal(t, 1, rowsIn(t, base), "the base row must be archived for the variants to mean anything") + + for _, kv := range parityBase { + field := kv[0] + kept, dropped := 0, 0 + for _, v := range valueVariants { + t.Run(field+"="+v, func(t *testing.T) { + if rowsIn(t, requireParity(t, archiveSchema, []parityInput{{topic: parityTopic, body: rowWith(field, v)}})) == 1 { + kept++ + } else { + dropped++ + } + }) + } + t.Logf("%-15s kept %2d of %d values, dropped %2d", field, kept, len(valueVariants), dropped) + if field != "exchange" && field != "day" { + require.Positive(t, dropped, "%s: no value was rejected, so the variants test nothing of its coercion", field) + } + require.Positive(t, kept, "%s: every value was rejected", field) + } +} + +func TestJSONParquetParityMessages(t *testing.T) { + valid := string(rowWith("", "")) + for name, body := range map[string]string{ + "empty": ``, + "whitespace": " \n\t ", + "not json": `not json`, + "truncated": valid[:len(valid)-1], + "array": `[` + valid + `]`, + "number": `5`, + "string": `"x"`, + "null": `null`, + "two documents": valid + valid, + "two documents spaced": valid + " " + valid, + "trailing garbage": valid + " x", + "trailing comma": strings.TrimSuffix(valid, "}") + ",}", + "surrounding whitespace": "\n " + valid + " \r\n", + "bom": "\xef\xbb\xbf" + valid, + "duplicate field": strings.TrimSuffix(valid, "}") + `,"price":2.5}`, + "duplicate null field": strings.TrimSuffix(valid, "}") + `,"price":null}`, + "duplicate required": `{"symbol":null,` + valid[1:], + "extra fields": strings.TrimSuffix(valid, "}") + `,"extra":{"deep":[1,2,{"x":"y"}]},"more":"é"}`, + "invalid utf8 in extra": strings.TrimSuffix(valid, "}") + ",\"extra\":\"\xff\"}", + "invalid utf8 in name": strings.TrimSuffix(valid, "}") + ",\"\xff\":1}", + "invalid json in extra": strings.TrimSuffix(valid, "}") + `,"extra":[1,}`, + "escaped field name": strings.Replace(valid, `"symbol"`, `"symbol"`, 1), + "deep nesting in extra": strings.TrimSuffix(valid, "}") + `,"extra":` + strings.Repeat("[", 5000) + strings.Repeat("]", 5000) + `}`, + "too deep in extra": strings.TrimSuffix(valid, "}") + `,"extra":` + strings.Repeat("[", 20000) + strings.Repeat("]", 20000) + `}`, + "long string": strings.Replace(valid, `"perpetual-BTC-USDT:USD"`, `"`+strings.Repeat("x", 1<<20)+`"`, 1), + "all fields missing": `{}`, + } { + t.Run(name, func(t *testing.T) { + requireParity(t, archiveSchema, []parityInput{{topic: parityTopic, body: []byte(body)}}) + }) + } +} + +func TestJSONParquetParityMetadata(t *testing.T) { + body := rowWith("", "") + for name, topic := range map[string]string{ + "no topic": "", + "short topic": "metric.v1", + "exactly 3 parts": "a.b.c", + "empty segment": "metric.v1..15s", + "unicode": "metric.v1.交易所.15s", + } { + t.Run(name, func(t *testing.T) { + requireParity(t, archiveSchema, []parityInput{{topic: topic, body: body}}) + }) + } +} + +func TestJSONParquetParityPartitions(t *testing.T) { + // the day of a row's time, across midnight, before the epoch and far ahead, + // and the path fields as every kind + var in []parityInput + for _, tm := range []string{ + `1790640000000000`, `1790639999999999`, `1790640000000001`, `1790726399999999`, `1790726400000000`, + `0`, `-1`, `-1000000`, `-86400000001`, `253402300799999999`, `4102444800000000`, + `1790677155000000.5`, `1.790677155e15`, `"1790677155000000"`, `1e400`, `null`, + } { + in = append(in, parityInput{topic: parityTopic, body: rowWith("time", tm)}) + } + for _, v := range []string{`15`, `true`, `null`, `{"a":1}`, `[1]`, `"1m"`, `"a/b"`, `"%s"`, `""`} { + in = append(in, parityInput{topic: parityTopic, body: rowWith("interval", v)}) + in = append(in, parityInput{topic: parityTopic, body: rowWith("timestamp_type", v)}) + } + for _, topic := range []string{"metric.v1.okx-perps.15s", "metric.v1.top-3.15s", parityTopic} { + in = append(in, parityInput{topic: topic, body: rowWith("", "")}) + } + files := requireParity(t, archiveSchema, in) + require.Greater(t, len(files), 5, "the batch must straddle many partitions to be a test of them") + for _, m := range in { + requireParity(t, archiveSchema, []parityInput{m}) + } +} + +// randomRow draws a row as production would send one, and now and then +// corrupts a field or the whole message. +func randomRow(r *rand.Rand) parityInput { + topics := []string{"metric.v1.binance-futures.15s", "metric.v1.okx-perps.1m", "metric.v1.top-3.30s", "metric.v1.hyperliquid-perps.1d"} + fields := [][2]string{ + {"exchange", fmt.Sprint(r.IntN(40))}, + {"symbol", fmt.Sprintf("%q", []string{"perpetual-BTC-USDT:USD", "perpetual-我踏马来了-USDT:USD", "perpetual-xyz:PALLADIUM-USD1", ""}[r.IntN(4)])}, + {"interval", fmt.Sprintf("%q", []string{"15s", "30s", "1m", "1d"}[r.IntN(4)])}, + {"timestamp_type", fmt.Sprintf("%q", []string{"true", "false"}[r.IntN(2)])}, + {"day", `"2026-09-29"`}, + {"time", fmt.Sprint(1790640000000000 + r.Int64N(3*86400)*1000000 - 86400000000)}, + {"count", fmt.Sprint(r.Int32() - (1 << 30))}, + {"count_opt", []string{fmt.Sprint(r.IntN(100)), "null"}[r.IntN(2)]}, + {"ratio", fmt.Sprint(r.Float32() * 100)}, + {"ratio_opt", []string{fmt.Sprint(r.NormFloat64()), "null"}[r.IntN(2)]}, + {"price", fmt.Sprint(r.ExpFloat64() * 1e4)}, + {"price_opt", []string{fmt.Sprint(-r.Float64()), "null"}[r.IntN(2)]}, + } + if r.IntN(10) == 0 { + f := &fields[r.IntN(len(fields))] + f[1] = valueVariants[r.IntN(len(valueVariants))] + } + var parts []string + for _, f := range fields { + if f[1] != missing { + parts = append(parts, fmt.Sprintf("%q:%s", f[0], f[1])) + } + } + r.Shuffle(len(parts), func(i, j int) { parts[i], parts[j] = parts[j], parts[i] }) + body := "{" + strings.Join(parts, ",") + "}" + if r.IntN(50) == 0 { + body = body[:r.IntN(len(body))] + } + return parityInput{topic: topics[r.IntN(len(topics))], body: []byte(body)} +} + +func TestJSONParquetParityRandomBatches(t *testing.T) { + r := rand.New(rand.NewPCG(1, 2)) + rows, inputs, files := 0, 0, 0 + for range 300 { + in := make([]parityInput, 1+r.IntN(400)) + for j := range in { + in[j] = randomRow(r) + } + out := requireParity(t, archiveSchema, in) + inputs += len(in) + files += len(out) + rows += rowsIn(t, out) + } + t.Logf("%d messages: %d rows archived in %d files, %d dropped", inputs, rows, files, inputs-rows) + require.Greater(t, rows, inputs*8/10, "most generated rows are valid") + require.Less(t, rows, inputs, "some generated rows must be dropped") +} + +func TestJSONParquetParityFullBatch(t *testing.T) { + // a 10 MiB batch, the byte budget an archiver closes a batch at: many pages + // per column and several row flushes per partition + r := rand.New(rand.NewPCG(3, 4)) + var in []parityInput + for size := 0; size < 10<<20; { + m := randomRow(r) + size += len(m.body) + in = append(in, m) + } + files := requireParity(t, archiveSchema, in) + t.Logf("%d messages: %d rows in %d files", len(in), rowsIn(t, files), len(files)) + require.Greater(t, rowsIn(t, files), len(in)*8/10) +} + +func TestJSONParquetParityWithoutPartitionOrNaN(t *testing.T) { + s := archiveSchema + s.partition, s.legacyPartition, s.nanForNull = "", "", false + r := rand.New(rand.NewPCG(5, 6)) + for range 50 { + in := make([]parityInput, 1+r.IntN(200)) + for j := range in { + in[j] = randomRow(r) + } + requireParity(t, s, in) + } + for _, v := range []string{`null`, missing, `1`} { + requireParity(t, s, []parityInput{{topic: parityTopic, body: rowWith("price", v)}}) + } +} + +func FuzzJSONParquetParity(f *testing.F) { + f.Add(rowWith("", ""), parityTopic) + for _, v := range valueVariants { + f.Add(rowWith("price", v), parityTopic) + f.Add(rowWith("time", v), "metric.v1.okx-perps.1m") + } + f.Fuzz(func(t *testing.T, body []byte, topic string) { + // newline-separated messages, so the fuzzer also varies batches + var in []parityInput + for line := range strings.SplitSeq(string(body), "\n") { + in = append(in, parityInput{topic: topic, body: []byte(line)}) + } + requireParity(t, archiveSchema, in) + }) +} diff --git a/internal/impl/parquet/processor_json_encode_parity_test.go b/internal/impl/parquet/processor_json_encode_parity_test.go new file mode 100644 index 0000000000..95071ba6f9 --- /dev/null +++ b/internal/impl/parquet/processor_json_encode_parity_test.go @@ -0,0 +1,368 @@ +package parquet_test + +import ( + "bytes" + "context" + "fmt" + "math" + "sort" + "strings" + "sync" + "testing" + "time" + + "github.com/parquet-go/parquet-go" + "github.com/stretchr/testify/require" + + _ "github.com/warpstreamlabs/bento/internal/impl/parquet" + _ "github.com/warpstreamlabs/bento/public/components/pure" + "github.com/warpstreamlabs/bento/public/service" +) + +// The parity harness runs the pipeline json_parquet_encode replaces and +// json_parquet_encode itself over the same batches, each as a real Bento +// stream, and requires the same Parquet files, byte for byte, under the same +// partitions with the same metadata, and the same rows dropped. +// +// The replaced pipeline is the one a Bloblang mapping per column builds: the +// mapping coerces every column and throws on a row it cannot, a catch drops the +// row, group_by_value splits the batch by the partition path the mapping put +// in metadata, and parquet_encode writes each group. + +type parityColumn struct { + name string + typ string + optional bool + fromMeta string // the Bloblang of an interpolated column, "" to read the field +} + +type paritySchema struct { + columns []parityColumn + partition string // json_parquet_encode's partition path, "" for none + // legacyPartition is the Bloblang expression the replaced mapping writes to + // meta partition; it must mean the same as partition. + legacyPartition string + nanForNull bool +} + +// archiveSchema has every column type and optionality, an interpolated column +// and a time-formatted partition, as metric rows archived by exchange and day. +var archiveSchema = paritySchema{ + columns: []parityColumn{ + {name: "exchange", typ: "UTF8", fromMeta: `@kafka_topic.split(".").index(2)`}, + {name: "symbol", typ: "UTF8"}, + {name: "interval", typ: "UTF8"}, + {name: "timestamp_type", typ: "UTF8"}, + {name: "day", typ: "UTF8"}, + {name: "time", typ: "INT64"}, + {name: "count", typ: "INT32"}, + {name: "count_opt", typ: "INT32", optional: true}, + {name: "ratio", typ: "FLOAT"}, + {name: "ratio_opt", typ: "FLOAT", optional: true}, + {name: "price", typ: "DOUBLE"}, + {name: "price_opt", typ: "DOUBLE", optional: true}, + }, + partition: "timestamp={timestamp_type}/{interval}/exchange={exchange}/{time|year=2006/month=01/day=02}", + legacyPartition: `"timestamp=%s/%s/exchange=%s/%s".format(this.timestamp_type, this.interval, root.exchange, (this.time / 1000000).ts_format("year=2006/month=01/day=02", "UTC"))`, + nanForNull: true, +} + +func (s paritySchema) schemaYAML() string { + var b strings.Builder + for _, c := range s.columns { + fmt.Fprintf(&b, " - { name: %s, type: %s, optional: %v }\n", c.name, c.typ, c.optional) + } + return b.String() +} + +// legacyMapping is the Bloblang the replaced pipeline coerces rows with. +func (s paritySchema) legacyMapping() string { + var required, checks []string + for _, c := range s.columns { + if c.fromMeta != "" { + continue + } + if !c.optional { + required = append(required, fmt.Sprintf("%q", c.name)) + } + integer := c.typ == "INT32" || c.typ == "INT64" + switch { + case c.typ == "UTF8": + checks = append(checks, fmt.Sprintf("root.%[1]s = this.%[1]s.not_null().string()", c.name)) + case integer && c.optional: + checks = append(checks, fmt.Sprintf("root.%[1]s = if this.%[1]s != null { this.%[1]s.%[2]s() } else { null }", c.name, strings.ToLower(c.typ))) + case integer: + checks = append(checks, fmt.Sprintf("root.%[1]s = this.%[1]s.%[2]s()", c.name, strings.ToLower(c.typ))) + case c.optional: + checks = append(checks, fmt.Sprintf("root.%[1]s = if this.%[1]s != null { this.%[1]s.number() } else { null }", c.name)) + case s.nanForNull: + checks = append(checks, fmt.Sprintf("root.%[1]s = this.%[1]s.or($nan).number()", c.name)) + default: + checks = append(checks, fmt.Sprintf("root.%[1]s = this.%[1]s.not_null().number()", c.name)) + } + } + lines := []string{ + "root = this", + "let missing = [" + strings.Join(required, ", ") + "].filter(c -> !this.exists(c))", + `root = if $missing.length() > 0 { throw("missing columns: " + $missing.join(", ")) } else { this }`, + `let nan = "NaN".number()`, + } + lines = append(lines, checks...) + for _, c := range s.columns { + if c.fromMeta != "" { + lines = append(lines, fmt.Sprintf("root.%s = %s", c.name, c.fromMeta)) + } + } + if s.legacyPartition != "" { + lines = append(lines, "meta partition = "+s.legacyPartition) + } + return strings.Join(lines, "\n") +} + +func (s paritySchema) legacyProcessors() string { + var b strings.Builder + b.WriteString("pipeline:\n processors:\n - mapping: |\n") + for l := range strings.SplitSeq(s.legacyMapping(), "\n") { + b.WriteString(" " + l + "\n") + } + b.WriteString(` - catch: + - mapping: root = deleted() +`) + if s.partition != "" { + b.WriteString(` - group_by_value: + value: '${! meta("partition") }' +`) + } + b.WriteString(" - parquet_encode:\n default_compression: zstd\n schema:\n") + b.WriteString(strings.ReplaceAll(s.schemaYAML(), " - ", " - ")) + return b.String() +} + +func (s paritySchema) newProcessors() string { + var b strings.Builder + fmt.Fprintf(&b, "pipeline:\n processors:\n - json_parquet_encode:\n default_compression: zstd\n nan_for_null: %v\n schema:\n", s.nanForNull) + b.WriteString(strings.ReplaceAll(s.schemaYAML(), " - ", " - ")) + b.WriteString(" columns:\n") + for _, c := range s.columns { + if c.fromMeta != "" { + fmt.Fprintf(&b, " - { name: %s, value: '${! %s }' }\n", c.name, c.fromMeta) + } + } + if s.partition != "" { + fmt.Fprintf(&b, " partition:\n path: '%s'\n time_unit: us\n", s.partition) + } + return b.String() +} + +//------------------------------------------------------------------------------ + +// parityStream is a running stream fed one batch at a time; each send returns +// once the batch is acknowledged, with what reached the output. +type parityStream struct { + send service.MessageBatchHandlerFunc + mu sync.Mutex + out []service.MessageBatch + stopped chan error // receives Run's error should the stream stop +} + +func startParityStream(tb testing.TB, processors string) *parityStream { + tb.Helper() + b := service.NewStreamBuilder() + require.NoError(tb, b.SetLoggerYAML("level: none")) + // processors is a pipeline section: each list item is added as a processor + body := strings.TrimPrefix(processors, "pipeline:\n processors:\n") + for _, item := range strings.Split("\n"+body, "\n - ")[1:] { + var lines []string + for l := range strings.SplitSeq(item, "\n") { + lines = append(lines, strings.TrimPrefix(l, " ")) + } + require.NoError(tb, b.AddProcessorYAML(strings.Join(lines, "\n")), item) + } + send, err := b.AddBatchProducerFunc() + require.NoError(tb, err) + ps := &parityStream{send: send} + require.NoError(tb, b.AddBatchConsumerFunc(func(_ context.Context, batch service.MessageBatch) error { + ps.mu.Lock() + ps.out = append(ps.out, batch.Copy()) + ps.mu.Unlock() + return nil + })) + stream, err := b.Build() + require.NoError(tb, err, processors) + ps.stopped = make(chan error, 1) + go func() { ps.stopped <- stream.Run(context.Background()) }() + return ps +} + +type parityInput struct { + topic string + body []byte +} + +type parityFile struct { + partition string + data []byte + meta map[string]any +} + +func (ps *parityStream) run(tb testing.TB, in []parityInput) []parityFile { + tb.Helper() + ps.mu.Lock() + ps.out = nil + ps.mu.Unlock() + batch := make(service.MessageBatch, len(in)) + for i, m := range in { + msg := service.NewMessage(append([]byte(nil), m.body...)) + if m.topic != "" { + msg.MetaSetMut("kafka_topic", m.topic) + } + msg.MetaSetMut("kafka_offset", i) + batch[i] = msg + } + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + sent := make(chan error, 1) + go func() { sent <- ps.send(ctx, batch) }() + select { + case err := <-sent: + require.NoError(tb, err) + case err := <-ps.stopped: + ps.stopped <- err + tb.Fatalf("the stream stopped: %v", err) + } + ps.mu.Lock() + defer ps.mu.Unlock() + var files []parityFile + for _, b := range ps.out { + for _, m := range b { + data, err := m.AsBytes() + require.NoError(tb, err) + meta := map[string]any{} + require.NoError(tb, m.MetaWalkMut(func(k string, v any) error { meta[k] = v; return nil })) + p, _ := m.MetaGet("partition") + files = append(files, parityFile{partition: p, data: append([]byte(nil), data...), meta: meta}) + } + } + return files +} + +type parityPair struct{ legacy, next *parityStream } + +var ( + parityPairsMu sync.Mutex + parityPairs = map[string]*parityPair{} +) + +// pairFor starts (once per schema) the replaced pipeline and json_parquet_encode. +func pairFor(tb testing.TB, s paritySchema) *parityPair { + tb.Helper() + key := s.legacyProcessors() + s.newProcessors() + parityPairsMu.Lock() + defer parityPairsMu.Unlock() + if p, ok := parityPairs[key]; ok { + return p + } + p := &parityPair{legacy: startParityStream(tb, s.legacyProcessors()), next: startParityStream(tb, s.newProcessors())} + parityPairs[key] = p + return p +} + +// requireParity runs a batch through both and fails on any difference. +func requireParity(tb testing.TB, s paritySchema, in []parityInput) []parityFile { + tb.Helper() + p := pairFor(tb, s) + legacy := p.legacy.run(tb, in) + next := p.next.run(tb, in) + describe := func(fs []parityFile) string { + var parts []string + for _, f := range fs { + parts = append(parts, fmt.Sprintf("%q (%d bytes, %d rows)", f.partition, len(f.data), len(parityRows(tb, f.data)))) + } + return strings.Join(parts, ", ") + } + if len(legacy) != len(next) { + tb.Fatalf("legacy wrote %d files [%s], json_parquet_encode %d [%s]\ninput: %s", len(legacy), describe(legacy), len(next), describe(next), describeInput(in)) + } + for i := range legacy { + l, n := legacy[i], next[i] + if l.partition != n.partition { + tb.Fatalf("file %d: partition %q, json_parquet_encode %q\ninput: %s", i, l.partition, n.partition, describeInput(in)) + } + if !bytes.Equal(l.data, n.data) { + lr, nr := parityRows(tb, l.data), parityRows(tb, n.data) + tb.Fatalf("file %d (%q) differs:\n%s\ninput: %s", i, l.partition, diffRows(lr, nr), describeInput(in)) + } + if fmt.Sprint(sortedMeta(l.meta)) != fmt.Sprint(sortedMeta(n.meta)) { + tb.Fatalf("file %d (%q): metadata %v, json_parquet_encode %v", i, l.partition, sortedMeta(l.meta), sortedMeta(n.meta)) + } + } + return next +} + +func describeInput(in []parityInput) string { + var parts []string + for _, m := range in { + parts = append(parts, fmt.Sprintf("[%s] %q", m.topic, m.body)) + } + if len(parts) > 12 { + parts = append(parts[:12], fmt.Sprintf("... %d more", len(parts)-12)) + } + return strings.Join(parts, "\n ") +} + +func sortedMeta(m map[string]any) []string { + var out []string + for k, v := range m { + out = append(out, fmt.Sprintf("%s=%v", k, v)) + } + sort.Strings(out) + return out +} + +// parityRows decodes a Parquet file into printable rows, float bits included. +func parityRows(tb testing.TB, data []byte) []string { + tb.Helper() + f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data))) + require.NoError(tb, err) + var rows []string + for _, rg := range f.RowGroups() { + r := rg.Rows() + buf := make([]parquet.Row, 64) + for { + n, err := r.ReadRows(buf) + for _, row := range buf[:n] { + var vals []string + for _, v := range row { + switch { + case v.IsNull(): + vals = append(vals, "NULL") + case v.Kind() == parquet.Float: + vals = append(vals, fmt.Sprintf("f32:%v(%08x)", v.Float(), math.Float32bits(v.Float()))) + case v.Kind() == parquet.Double: + vals = append(vals, fmt.Sprintf("f64:%v(%016x)", v.Double(), math.Float64bits(v.Double()))) + default: + vals = append(vals, v.String()) + } + } + rows = append(rows, strings.Join(vals, " | ")) + } + if err != nil { + break + } + } + _ = r.Close() + } + return rows +} + +func diffRows(legacy, next []string) string { + if len(legacy) != len(next) { + return fmt.Sprintf("legacy %d rows, json_parquet_encode %d rows\nlegacy: %v\nnew: %v", len(legacy), len(next), legacy, next) + } + for i := range legacy { + if legacy[i] != next[i] { + return fmt.Sprintf("row %d:\nlegacy: %s\nnew: %s", i, legacy[i], next[i]) + } + } + return "same rows, different bytes (encoding or file metadata)" +} diff --git a/internal/impl/parquet/processor_json_encode_test.go b/internal/impl/parquet/processor_json_encode_test.go new file mode 100644 index 0000000000..e2bf0593c9 --- /dev/null +++ b/internal/impl/parquet/processor_json_encode_test.go @@ -0,0 +1,179 @@ +package parquet + +import ( + "bytes" + "context" + "math" + "testing" + + "github.com/parquet-go/parquet-go" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/warpstreamlabs/bento/public/service" +) + +func newTestJSONEncoder(t *testing.T, conf string) (*jsonParquetEncoder, error) { + t.Helper() + parsed, err := jsonParquetEncodeSpec().ParseYAML(conf, nil) + require.NoError(t, err) + return newJSONParquetEncoder(parsed, service.MockResources()) +} + +func TestJSONParquetEncodeRejectsWhatItCannotWrite(t *testing.T) { + // A config json_parquet_encode would write differently from parquet_encode, + // or not at all, must fail at start rather than on the first batch. + for name, conf := range map[string]string{ + "no columns": `schema: []`, + "nested column": `schema: [ { name: a, type: STRUCT, fields: [ { name: b, type: UTF8 } ] } ]`, + "list column": `schema: [ { name: a, type: LIST, fields: [ { name: element, type: UTF8 } ] } ]`, + "unsupported type": `schema: [ { name: a, type: BOOLEAN } ]`, + "byte array": `schema: [ { name: a, type: BYTE_ARRAY } ]`, + "repeated": `schema: [ { name: a, type: INT64, repeated: true } ]`, + "duplicate column": `schema: [ { name: a, type: INT64 }, { name: a, type: UTF8 } ]`, + "unknown override": "schema: [ { name: a, type: UTF8 } ]\ncolumns: [ { name: b, value: x } ]", + "override not UTF8": "schema: [ { name: a, type: INT64 } ]\ncolumns: [ { name: a, value: '1' } ]", + "unknown placeholder": "schema: [ { name: a, type: UTF8 } ]\npartition: { path: 'x={b}' }", + "unclosed placeholder": "schema: [ { name: a, type: UTF8 } ]\npartition: { path: 'x={a' }", + "empty time layout": "schema: [ { name: a, type: INT64 } ]\npartition: { path: 'x={a|}' }", + "time of interpolated": "schema: [ { name: a, type: UTF8 } ]\ncolumns: [ { name: a, value: x } ]\npartition: { path: '{a|2006}' }", + "unknown time unit": "schema: [ { name: a, type: INT64 } ]\npartition: { path: '{a|2006}', time_unit: days }", + } { + t.Run(name, func(t *testing.T) { + _, err := newTestJSONEncoder(t, conf) + assert.Error(t, err) + }) + } +} + +func readRows(t *testing.T, data []byte) []parquet.Row { + t.Helper() + f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data))) + require.NoError(t, err) + var rows []parquet.Row + for _, rg := range f.RowGroups() { + r := rg.Rows() + buf := make([]parquet.Row, 16) + for { + n, err := r.ReadRows(buf) + for _, row := range buf[:n] { + rows = append(rows, row.Clone()) + } + if err != nil { + break + } + } + require.NoError(t, r.Close()) + } + return rows +} + +func TestJSONParquetEncodeDropsInvalidRows(t *testing.T) { + // Dropped rows are acknowledged with the batch and never archived; the rest + // of the batch is written. + parsed, err := jsonParquetEncodeSpec().ParseYAML(` +schema: + - { name: id, type: INT64 } + - { name: v, type: DOUBLE, optional: true } +`, nil) + require.NoError(t, err) + res := service.MockResources() + e, err := newJSONParquetEncoder(parsed, res) + require.NoError(t, err) + + out, err := e.ProcessBatch(context.Background(), service.MessageBatch{ + service.NewMessage([]byte(`{"id":1,"v":null}`)), + service.NewMessage([]byte(`{"v":2}`)), // missing id + service.NewMessage([]byte(`{"id":1.5}`)), // not an integer + service.NewMessage([]byte(`not json`)), // not JSON + service.NewMessage([]byte(`{"id":2,"v":3}`)), + }) + require.NoError(t, err) + require.Len(t, out, 1) + require.Len(t, out[0], 1, "without a partition the batch is one file") + data, err := out[0][0].AsBytes() + require.NoError(t, err) + rows := readRows(t, data) + require.Len(t, rows, 2) + assert.Equal(t, int64(1), rows[0][0].Int64()) + assert.True(t, rows[0][1].IsNull()) + assert.Equal(t, int64(2), rows[1][0].Int64()) + assert.Equal(t, 3.0, rows[1][1].Double()) +} + +func TestJSONParquetEncodeAllDroppedWritesNothing(t *testing.T) { + e, err := newTestJSONEncoder(t, `schema: [ { name: id, type: INT64 } ]`) + require.NoError(t, err) + out, err := e.ProcessBatch(context.Background(), service.MessageBatch{service.NewMessage([]byte(`{}`))}) + require.NoError(t, err) + assert.Nil(t, out, "an empty Parquet file would be uploaded as an object with no rows") + out, err = e.ProcessBatch(context.Background(), service.MessageBatch{}) + require.NoError(t, err) + assert.Nil(t, out) +} + +func TestJSONParquetEncodePartitionsKeepFirstMessageMetadata(t *testing.T) { + e, err := newTestJSONEncoder(t, ` +schema: + - { name: ex, type: UTF8 } + - { name: t, type: INT64 } +columns: + - { name: ex, value: '${! @topic }' } +partition: + path: 'ex={ex}/{t|2006-01-02}' + time_unit: s + metadata_key: part +`) + require.NoError(t, err) + msg := func(topic, body string, offset int) *service.Message { + m := service.NewMessage([]byte(body)) + m.MetaSetMut("topic", topic) + m.MetaSetMut("offset", offset) + return m + } + out, err := e.ProcessBatch(context.Background(), service.MessageBatch{ + msg("a", `{"t":86399}`, 1), // 1970-01-01 + msg("b", `{"t":86400}`, 2), // another exchange + msg("a", `{"t":86400}`, 3), // the next day + msg("a", `{"t":0}`, 4), // back to the first partition + }) + require.NoError(t, err) + require.Len(t, out, 1) + var parts []string + for _, m := range out[0] { + p, _ := m.MetaGet("part") + off, _ := m.MetaGetMut("offset") + parts = append(parts, p) + data, err := m.AsBytes() + require.NoError(t, err) + switch p { + case "ex=a/1970-01-01": + assert.Equal(t, 1, off, "a file carries its first message's metadata") + assert.Len(t, readRows(t, data), 2) + default: + assert.Len(t, readRows(t, data), 1) + } + } + assert.Equal(t, []string{"ex=a/1970-01-01", "ex=b/1970-01-02", "ex=a/1970-01-02"}, parts, "files in order of first appearance") +} + +func TestJSONParquetEncodeOptionalStringIsNull(t *testing.T) { + // No Bloblang the archivers generate has an optional UTF8 column; here an + // optional column of any type reads null (or its absence) as NULL. + e, err := newTestJSONEncoder(t, `schema: [ { name: s, type: UTF8, optional: true }, { name: n, type: FLOAT } ]`) + require.NoError(t, err) + out, err := e.ProcessBatch(context.Background(), service.MessageBatch{ + service.NewMessage([]byte(`{"s":null,"n":1}`)), + service.NewMessage([]byte(`{"n":2}`)), + service.NewMessage([]byte(`{"s":"x","n":null}`)), // null in a required float, nan_for_null off + }) + require.NoError(t, err) + data, err := out[0][0].AsBytes() + require.NoError(t, err) + rows := readRows(t, data) + require.Len(t, rows, 2) + assert.True(t, rows[0][0].IsNull()) + assert.True(t, rows[1][0].IsNull()) + assert.Equal(t, float32(2), rows[1][1].Float()) + assert.False(t, math.IsNaN(float64(rows[1][1].Float()))) +} diff --git a/website/docs/components/inputs/kafka_franz.md b/website/docs/components/inputs/kafka_franz.md index 1be209a82b..366cf32dc0 100644 --- a/website/docs/components/inputs/kafka_franz.md +++ b/website/docs/components/inputs/kafka_franz.md @@ -76,6 +76,7 @@ input: sasl: [] # No default (optional) multi_header: false add_record_metadata: true + copy_record_values: false batching: count: 0 byte_size: 0 @@ -779,6 +780,14 @@ Add the `kafka_*` metadata fields and record headers to each message. Disabling Type: `bool` Default: `true` +### `copy_record_values` + +Copy each record's value into its message instead of referencing the fetch response it arrived in. A record's key, value and header values are slices of that response, which holds every record fetched from the broker in the same request, so a message held for long (e.g. in a large output batch) otherwise keeps the whole response in memory. Costs one allocation per record. + + +Type: `bool` +Default: `false` + ### `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/json_parquet_encode.md b/website/docs/components/processors/json_parquet_encode.md new file mode 100644 index 0000000000..2b35466e25 --- /dev/null +++ b/website/docs/components/processors/json_parquet_encode.md @@ -0,0 +1,277 @@ +--- +title: json_parquet_encode +slug: json_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. +::: +Encodes a batch of flat JSON objects straight into Parquet, row by row, without holding the batch as structured messages. + + + + + + +```yml +# Common config fields, showing default values +label: "" +json_parquet_encode: + schema: [] # No default (required) + columns: [] + nan_for_null: false + partition: + path: day={time|2006-01-02}/{symbol} # No default (required) + time_unit: us + metadata_key: partition + default_compression: uncompressed +``` + + + + +```yml +# All config fields, showing default values +label: "" +json_parquet_encode: + schema: [] # No default (required) + columns: [] + nan_for_null: false + partition: + path: day={time|2006-01-02}/{symbol} # No default (required) + time_unit: us + metadata_key: partition + default_compression: uncompressed + default_encoding: DELTA_LENGTH_BYTE_ARRAY +``` + + + + +Produces the same Parquet as chaining a `mapping` that coerces each column, a `catch` that drops the rows it rejects, `group_by_value` on a partition path and `parquet_encode`, at a fraction of the memory: a batch is held only as its raw bytes until it is processed, each message is then parsed once into a Parquet row, and the files are encoded one at a time. + +The schema is `parquet_encode`'s, restricted to flat `UTF8`, `INT32`, `INT64`, `FLOAT` and `DOUBLE` columns, and the file is written by the same encoder. A column reads the message field of its name and coerces it as Bloblang would: + +- `UTF8`: `.not_null().string()`, so a number keeps its literal text; +- `INT32`, `INT64`: `.int32()`, `.int64()` (a fraction, an exponent or an out-of-range value is rejected; a string is parsed); +- `FLOAT`, `DOUBLE`: `.number()`, narrowed to 32 bits for `FLOAT`. + +A required column the message lacks rejects it. A `null` in an optional column is written as NULL, and in a required float column as NaN when `nan_for_null` is set; otherwise it rejects the message. A column may instead take an interpolated `value`, evaluated per message against its metadata only. + +A message that is not a single JSON object, or that any column rejects, is dropped: it is logged, counted in `json_parquet_encode_dropped` and acknowledged with the batch. + +When `partition` is set the batch is split into one Parquet file per distinct partition path, in order of first appearance, and the path is written to the configured metadata key. `{column}` in the path is replaced by the column's value as Bloblang's `format` would print the field, and `{column|layout}` by the column's number, read in `time_unit`, as a UTC time in the Go layout. Every output message carries the metadata of the first message of its partition. + +## Examples + + + + + +Batches rows at the output and uploads one Parquet file per exchange and UTC day of a microsecond timestamp. + +```yaml +output: + aws_s3: + bucket: metrics + path: 'ohlcv/${! meta("partition") }/${! uuid_v4() }.parquet' + batching: + byte_size: 10485760 + period: 30s + processors: + - json_parquet_encode: + default_compression: zstd + nan_for_null: true + schema: + - { name: exchange, type: UTF8 } + - { name: symbol, type: UTF8 } + - { name: time, type: INT64 } + - { name: close, type: DOUBLE } + - { name: volume, type: DOUBLE, optional: true } + columns: + - { name: exchange, value: '${! @kafka_topic.split(".").index(2) }' } + partition: + path: 'exchange={exchange}/{time|year=2006/month=01/day=02}' + time_unit: us +``` + + + + +## Fields + +### `schema` + +Parquet schema. + + +Type: `array` + +### `schema[].name` + +The name of the column. + + +Type: `string` + +### `schema[].type` + +The type of the column, only applicable for leaf columns with no child fields. STRUCT represents nested objects with defined field schemas. MAP supports only string keys, but can support values of all types. Some logical types can be specified here such as UTF8. + + +Type: `string` +Options: `BOOLEAN`, `INT8`, `INT16`, `INT32`, `INT64`, `DECIMAL64`, `DECIMAL32`, `FLOAT`, `DOUBLE`, `BYTE_ARRAY`, `UTF8`, `MAP`, `LIST`, `STRUCT`. + +### `schema[].decimal_precision` + +Precision to use for DECIMAL32/DECIMAL64 type + + +Type: `int` +Default: `0` + +### `schema[].decimal_scale` + +Scale to use for DECIMAL32/DECIMAL64 type + + +Type: `int` +Default: `0` + +### `schema[].repeated` + +Whether the field is repeated. + + +Type: `bool` +Default: `false` + +### `schema[].optional` + +Whether the field is optional. + + +Type: `bool` +Default: `false` + +### `schema[].fields` + +A list of child fields. + + +Type: `array` + +```yml +# Examples + +fields: + - name: foo + type: INT64 + - name: bar + type: BYTE_ARRAY +``` + +### `columns` + +Columns set from an interpolation instead of the message field of their name. + + +Type: `array` +Default: `[]` + +### `columns[].name` + +The schema column this sets. + + +Type: `string` + +### `columns[].value` + +The column's value, from the message's metadata. It must not reference the message content. +This field supports [interpolation functions](/docs/configuration/interpolation#bloblang-queries). + + +Type: `string` + +### `nan_for_null` + +Write a `null` in a required `FLOAT` or `DOUBLE` column as NaN instead of dropping the message. + + +Type: `bool` +Default: `false` + +### `partition` + +Split the batch into one Parquet file per partition path. + + +Type: `object` + +### `partition.path` + +The partition path, with `{column}` and `{column|layout}` placeholders. + + +Type: `string` + +```yml +# Examples + +path: day={time|2006-01-02}/{symbol} +``` + +### `partition.time_unit` + +The unit of numbers formatted as a time. + + +Type: `string` +Default: `"us"` +Options: `s`, `ms`, `us`, `ns`. + +### `partition.metadata_key` + +The metadata field that receives the partition path. + + +Type: `string` +Default: `"partition"` + +### `default_compression` + +The default compression type to use for fields. + + +Type: `string` +Default: `"uncompressed"` +Options: `uncompressed`, `snappy`, `gzip`, `brotli`, `zstd`, `lz4raw`. + +### `default_encoding` + +The default encoding type to use for fields. + + +Type: `string` +Default: `"DELTA_LENGTH_BYTE_ARRAY"` +Options: `DELTA_LENGTH_BYTE_ARRAY`, `PLAIN`, `RLE_DICTIONARY`. + +