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
116 changes: 116 additions & 0 deletions .github/workflows/aperiodic_test.yml
Original file line number Diff line number Diff line change
@@ -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 }}
20 changes: 19 additions & 1 deletion internal/impl/kafka/input_kafka_franz.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package kafka

import (
"bytes"
"context"
"crypto/tls"
"errors"
Expand Down Expand Up @@ -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()).
Expand Down Expand Up @@ -229,6 +231,7 @@ type franzKafkaReader struct {
regexPattern bool
multiHeader bool
addMetadata bool
copyValues bool
batchPolicy service.BatchPolicy

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