Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions .github/workflows/aperiodic_image.yml
Original file line number Diff line number Diff line change
@@ -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
3 changes: 2 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
10 changes: 10 additions & 0 deletions internal/impl/kafka/input_kafka_franz.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()).
Expand Down Expand Up @@ -227,6 +228,7 @@ type franzKafkaReader struct {
commitPeriod time.Duration
regexPattern bool
multiHeader bool
addMetadata bool
batchPolicy service.BatchPolicy

reconnectOnUnknownTopic bool
Expand Down Expand Up @@ -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
}
Expand All @@ -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))
}
Expand Down
41 changes: 41 additions & 0 deletions internal/impl/kafka/input_kafka_franz_metadata_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading