Skip to content
Merged
Show file tree
Hide file tree
Changes from 7 commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
121d72d
[receiver/kafkareceiver] move receiver.kafkareceiver.UseFranzGo featu…
paulojmdias Aug 21, 2025
6fb2ecf
fix: make generate
paulojmdias Aug 21, 2025
da5da50
Merge branch 'main' into feat/42155
paulojmdias Aug 21, 2025
a002b19
fix: Update kafka_receiver_test.go
paulojmdias Aug 22, 2025
f1e6eca
fiz: context.Background()
paulojmdias Aug 22, 2025
09014ff
chore: //nolint:usetesting
paulojmdias Aug 22, 2025
a4131e8
Merge branch 'main' into feat/42155
paulojmdias Aug 22, 2025
f56fea2
chore: update readme.md
paulojmdias Aug 25, 2025
be8ae76
feat: test and consumer improvements
paulojmdias Aug 26, 2025
060d629
Merge branch 'main' into feat/42155
paulojmdias Aug 26, 2025
54afbbc
fix: improve tests and revert shutdown ignore cancellation
paulojmdias Aug 27, 2025
e5ddaf7
Merge branch 'main' into feat/42155
paulojmdias Aug 27, 2025
7b7c607
feat: add support for profiles in tests
paulojmdias Aug 27, 2025
acd16ed
Merge branch 'main' into feat/42155
paulojmdias Sep 2, 2025
0ef7943
Merge branch 'main' into feat/42155
paulojmdias Sep 25, 2025
4d81c8d
Merge branch 'main' into feat/42155
paulojmdias Sep 25, 2025
566350f
fix: fix tests
paulojmdias Sep 25, 2025
aacea05
fix: improvements on signals
paulojmdias Sep 25, 2025
2c468ed
fix: improve test for CI stability
paulojmdias Sep 25, 2025
bda6657
Merge branch 'main' into feat/42155
paulojmdias Sep 25, 2025
f9a0d04
Merge branch 'main' into feat/42155
paulojmdias Sep 25, 2025
c33bdd0
Merge branch 'main' into feat/42155
paulojmdias Sep 25, 2025
de4d9a1
Update consumer_franz_test.go
paulojmdias Sep 25, 2025
0a21c09
Merge branch 'main' into feat/42155
paulojmdias Sep 26, 2025
de5a034
Merge branch 'main' into feat/42155
paulojmdias Sep 26, 2025
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
27 changes: 27 additions & 0 deletions .chloggen/receiver_kafka_franz_go_beta.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
# Use this changelog template to create an entry for release notes.

# One of 'breaking', 'deprecation', 'new_component', 'enhancement', 'bug_fix'
change_type: "enhancement"

# The name of the component, or a single word describing the area of concern, (e.g. filelogreceiver)
component: receiver/kafkareceiver

# A brief description of the change. Surround your text with quotes ("") if it needs to start with a backtick (`).
note: "Use franz-go client for Kafka receiver as default, promoting the receiver.kafkareceiver.UseFranzGo feature gate to Beta."

# Mandatory: One or more tracking issues related to the change. You can use the PR number here if no issue exists.
issues: [42155]

# (Optional) One or more lines of additional information to render under the primary note.
# These lines will be padded with 2 spaces and then inserted directly into the document.
# Use pipe (|) for multiline entries.
subtext:

# If your change doesn't affect end users or the exported elements of any package,
# you should instead start your pull request title with [chore] or use the "Skip Changelog" label.
# Optional: The change log or logs in which this entry should be included.
# e.g. '[user]' or '[user, api]'
# Include 'user' if the change is relevant to end users.
# Include 'api' if there is a change to a library API.
# Default: '[user]'
change_logs: [user]
10 changes: 5 additions & 5 deletions receiver/kafkareceiver/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ If used in conjunction with the `kafkaexporter` configured with `include_metadat
## Getting Started

> [!NOTE]
> You can opt-in to use [`franz-go`](https://github.com/twmb/franz-go) client by enabling the feature gate
> You can opt-out to use [`franz-go`](https://github.com/twmb/franz-go) client by disabling the feature gate
> `receiver.kafkareceiver.UseFranzGo` when you run the OpenTelemetry Collector. See the following page
> for more details: [Feature Gates](https://github.com/open-telemetry/opentelemetry-collector/tree/main/featuregate#controlling-gates)
>
Expand Down Expand Up @@ -109,10 +109,10 @@ The following settings can be optionally configured:
**Note: this can block the entire partition in case a message processing returns a permanent error**
- `header_extraction`:
- `extract_headers` (default = false): Allows user to attach header fields to resource attributes in otel pipeline
- `headers` (default = []): List of headers they'd like to extract from kafka record.
**Note: Matching pattern will be `exact`. Regexes are not supported as of now.**
- `headers` (default = []): List of headers they'd like to extract from kafka record.
**Note: Matching pattern will be `exact`. Regexes are not supported as of now.**
- `error_backoff`: [BackOff](https://github.com/open-telemetry/opentelemetry-collector/blob/v0.116.0/config/configretry/backoff.go#L27-L43) configuration in case of errors
- `enabled`: (default = false) Whether to enable backoff when next consumers return errors
- `enabled`: (default = false) Whether to enable backoff when next consumers return errors
- `initial_interval`: The time to wait after the first error before retrying
- `max_interval`: The upper bound on backoff interval between consecutive retries
- `multiplier`: The value multiplied by the backoff interval bounds
Expand Down Expand Up @@ -190,7 +190,7 @@ be configured to extract and attach specific headers as resource attributes. e.g
```yaml
receivers:
kafka:
header_extraction:
header_extraction:
extract_headers: true
headers: ["header1", "header2"]
```
Expand Down
4 changes: 2 additions & 2 deletions receiver/kafkareceiver/consumer_franz.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ const franzGoConsumerFeatureGateName = "receiver.kafkareceiver.UseFranzGo"
// the Kafka receiver will use the franz-go client, which is more performant and has
// better support for modern Kafka features.
var franzGoConsumerFeatureGate = featuregate.GlobalRegistry().MustRegister(
franzGoConsumerFeatureGateName, featuregate.StageAlpha,
franzGoConsumerFeatureGateName, featuregate.StageBeta,
featuregate.WithRegisterDescription("When enabled, the Kafka receiver will use the franz-go client to consume messages."),
featuregate.WithRegisterFromVersion("v0.129.0"),
)
Expand Down Expand Up @@ -351,7 +351,7 @@ func (c *franzConsumer) consume(ctx context.Context, size int) bool {

func (c *franzConsumer) Shutdown(ctx context.Context) error {
if !c.triggerShutdown() {
return errors.New("kafka consumer: consumer isn't running")
return nil
}

select {
Expand Down
8 changes: 6 additions & 2 deletions receiver/kafkareceiver/consumer_franz_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -238,10 +238,14 @@ func TestConsumerShutdownNotStarted(t *testing.T) {
c, err := newFranzKafkaConsumer(cfg, settings, []string{"test"}, nil)
require.NoError(t, err)

// `Shutdown` must be idempotent and succeed even if never started.
for i := 0; i < 2; i++ {
require.EqualError(t, c.Shutdown(t.Context()),
"kafka consumer: consumer isn't running")
require.NoError(t, c.Shutdown(t.Context()))
}

// Verify internal signal that there's nothing to shut down.
// (Same package, so we can call the unexported helper.)
require.False(t, c.triggerShutdown(), "triggerShutdown should indicate no-op when never started")
}

// TestRaceLostVsConsume verifies no data race occurs between concurrent
Expand Down
134 changes: 107 additions & 27 deletions receiver/kafkareceiver/factory_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,32 @@ import (
"github.com/open-telemetry/opentelemetry-collector-contrib/receiver/kafkareceiver/internal/metadata"
)

func encodingFromReceiver(t *testing.T, r any, section string) string {
t.Helper()
switch rc := r.(type) {
case *saramaConsumer:
switch section {
case "Traces":
return rc.config.Traces.Encoding
case "Metrics":
return rc.config.Metrics.Encoding
case "Logs":
return rc.config.Logs.Encoding
}
case *franzConsumer:
switch section {
case "Traces":
return rc.config.Traces.Encoding
case "Metrics":
return rc.config.Metrics.Encoding
case "Logs":
return rc.config.Logs.Encoding
}
}
t.Fatalf("unsupported receiver type %T or section %q", r, section)
return ""
}

func TestCreateDefaultConfig(t *testing.T) {
cfg := createDefaultConfig().(*Config)
assert.NotNil(t, cfg, "failed to create default config")
Expand All @@ -36,24 +62,42 @@ func TestCreateTraces(t *testing.T) {
func TestWithTracesUnmarshalers(t *testing.T) {
f := NewFactory()

t.Run("custom_encoding", func(t *testing.T) {
t.Run("custom_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig().(*Config)
cfg.Traces.Encoding = "custom"
receiver, err := f.CreateTraces(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Traces"))
})

t.Run("custom_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig().(*Config)
cfg.Traces.Encoding = "custom"
receiver, err := f.CreateTraces(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
tracesConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, "custom", tracesConsumer.config.Traces.Encoding)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Traces"))
})

t.Run("default_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig()
receiver, err := f.CreateTraces(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, defaultTracesEncoding, encodingFromReceiver(t, receiver, "Traces"))
})
t.Run("default_encoding", func(t *testing.T) {

t.Run("default_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig()
receiver, err := f.CreateTraces(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
tracesConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, defaultTracesEncoding, tracesConsumer.config.Traces.Encoding)
require.NoError(t, err)
assert.NotNil(t, receiver)
require.NotNil(t, receiver)
assert.Equal(t, defaultTracesEncoding, encodingFromReceiver(t, receiver, "Traces"))
})
}

Expand All @@ -70,24 +114,42 @@ func TestCreateMetrics(t *testing.T) {
func TestWithMetricsUnmarshalers(t *testing.T) {
f := NewFactory()

t.Run("custom_encoding", func(t *testing.T) {
t.Run("custom_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig().(*Config)
cfg.Metrics.Encoding = "custom"
receiver, err := f.CreateMetrics(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Metrics"))
})

t.Run("custom_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig().(*Config)
cfg.Metrics.Encoding = "custom"
receiver, err := f.CreateMetrics(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
metricsConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, "custom", metricsConsumer.config.Metrics.Encoding)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Metrics"))
})
t.Run("default_encoding", func(t *testing.T) {

t.Run("default_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig()
receiver, err := f.CreateMetrics(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, defaultMetricsEncoding, encodingFromReceiver(t, receiver, "Metrics"))
})

t.Run("default_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig()
receiver, err := f.CreateMetrics(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
metricsConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, defaultMetricsEncoding, metricsConsumer.config.Metrics.Encoding)
require.NoError(t, err)
assert.NotNil(t, receiver)
require.NotNil(t, receiver)
assert.Equal(t, defaultMetricsEncoding, encodingFromReceiver(t, receiver, "Metrics"))
})
}

Expand All @@ -104,23 +166,41 @@ func TestCreateLogs(t *testing.T) {
func TestWithLogsUnmarshalers(t *testing.T) {
f := NewFactory()

t.Run("custom_encoding", func(t *testing.T) {
t.Run("custom_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig().(*Config)
cfg.Logs.Encoding = "custom"
receiver, err := f.CreateLogs(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
logsConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, "custom", logsConsumer.config.Logs.Encoding)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Logs"))
})
t.Run("default_encoding", func(t *testing.T) {

t.Run("custom_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig().(*Config)
cfg.Logs.Encoding = "custom"
receiver, err := f.CreateLogs(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, "custom", encodingFromReceiver(t, receiver, "Logs"))
})

t.Run("default_encoding/sarama", func(t *testing.T) {
setFranzGo(t, false)
cfg := createDefaultConfig()
receiver, err := f.CreateLogs(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
logsConsumer, ok := receiver.(*saramaConsumer)
require.True(t, ok)
require.Equal(t, defaultLogsEncoding, logsConsumer.config.Logs.Encoding)
require.NoError(t, err)
assert.NotNil(t, receiver)
require.NotNil(t, receiver)
assert.Equal(t, defaultLogsEncoding, encodingFromReceiver(t, receiver, "Logs"))
})

t.Run("default_encoding/franzgo", func(t *testing.T) {
setFranzGo(t, true)
cfg := createDefaultConfig()
receiver, err := f.CreateLogs(t.Context(), receivertest.NewNopSettings(metadata.Type), cfg, nil)
require.NoError(t, err)
require.NotNil(t, receiver)
assert.Equal(t, defaultLogsEncoding, encodingFromReceiver(t, receiver, "Logs"))
})
}
Loading
Loading