A note for the community
- Please vote on this issue by adding a 👍 reaction to the original issue to help the community and maintainers prioritize this request
- If you are interested in working on this issue or have submitted a pull request, please leave a comment
Problem
Scenario:
A single node kafka broker, with 1 topic and one partition with 1 replica. Example docker-compose:
version: ‘3.7'
services:
zookeeper:
image: confluentinc/cp-zookeeper:latest
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
ports:
- 2181:2181
kafka:
image: confluentinc/cp-kafka:latest
depends_on:
- zookeeper
ports:
- 9092:9092
- 29092:29092
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
Vector source sending to the blackhole sink, example config:
[api]
enabled = true
[sources.kafka]
type = "kafka"
bootstrap_servers = "localhost:29092"
group_id = "aeohkh4k1j"
topics = ["conn"]
auto_offset_reset = "latest"
commit_interval_ms = 5000
metrics.topic_lag_metric = true
session_timeout_ms = 30000
[sinks.my_sink_id]
type = "blackhole"
inputs = [ "kafka" ]
I wrote a short python script to publish 150k messages per second. (just pip install confluent-kafka as a prereq):
import time
import json
from confluent_kafka import Producer
# Kafka configuration with increased buffer size and batch configuration
conf = {
'bootstrap.servers': 'localhost:29092',
'queue.buffering.max.messages': 1000000,
'queue.buffering.max.kbytes': 1048576,
'linger.ms': 10,
'batch.num.messages': 10000,
'compression.codec': 'gzip',
}
# Create a Kafka producer
producer = Producer(conf)
# Topic name
topic = 'conn'
# Define a static Zeek conn.log message
static_conn_log = {
"ts": "2023-10-05T12:34:56.789Z", # Fixed timestamp
"uid": "C1LZd81Vt3hZk3x5f2", # Fixed UID
"id.orig_h": "192.168.1.1",
"id.orig_p": 12345,
"id.resp_h": "192.168.1.2",
"id.resp_p": 80,
"proto": "tcp",
"duration": 1.234,
"orig_bytes": 500,
"resp_bytes": 1000,
"conn_state": "SF"
}
# Convert the static log to a JSON string
static_message = json.dumps(static_conn_log)
def produce_events(events_per_second):
"""Produce a specified number of events per second."""
event_count = 0
start_time = time.time()
while True:
current_time = time.time()
elapsed_time = current_time - start_time
# Calculate how many events should be produced in the elapsed time
expected_event_count = int(elapsed_time * events_per_second)
while event_count < expected_event_count:
try:
# Produce the static event
producer.produce(topic, key=None, value=static_message)
event_count += 1
except BufferError:
# If the local queue is full, wait briefly and retry
producer.poll(0.01)
# Poll to process background tasks
producer.poll(0)
# Sleep to avoid busy waiting
time.sleep(0.001)
try:
produce_events(100000) # 100,000 events per second
except KeyboardInterrupt:
pass
finally:
# Wait for any outstanding messages to be delivered
producer.flush()
I observed through vector top that the kafka consume side could only consume 100-200k messages at best, even though hardware resources were available. No amount of librdkafka_options tweaking could get better performance.

To prove better performance was available on the same hardware I wrote the following go program:
package main
import (
"context"
"errors"
"flag"
"github.com/IBM/sarama"
"log"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
)
// Sarama configuration options
var (
brokers = ""
version = ""
group = ""
topics = ""
assignor = ""
oldest = true
verbose = false
)
// Consumer represents a Sarama consumer group consumer
type Consumer struct {
ready chan bool
}
const (
batchSize = 1000 // Number of messages in a batch
batchTime = 100 * time.Millisecond // Maximum time to wait before sending a batch
)
func init() {
flag.StringVar(&brokers, "brokers", "", "Kafka bootstrap brokers to connect to, as a comma separated list")
flag.StringVar(&group, "group", "", "Kafka consumer group definition")
flag.StringVar(&version, "version", sarama.DefaultVersion.String(), "Kafka cluster version")
flag.StringVar(&topics, "topics", "", "Kafka topics to be consumed, as a comma separated list")
flag.StringVar(&assignor, "assignor", "range", "Consumer group partition assignment strategy (range, roundrobin, sticky)")
flag.BoolVar(&oldest, "oldest", true, "Kafka consumer consume initial offset from oldest")
flag.BoolVar(&verbose, "verbose", false, "Sarama logging")
flag.Parse()
if len(brokers) == 0 {
panic("no Kafka bootstrap brokers defined, please set the -brokers flag")
}
if len(topics) == 0 {
panic("no topics given to be consumed, please set the -topics flag")
}
if len(group) == 0 {
panic("no Kafka consumer group defined, please set the -group flag")
}
}
func main() {
keepRunning := true
log.Println("Starting a new Sarama consumer")
if verbose {
sarama.Logger = log.New(os.Stdout, "[sarama] ", log.LstdFlags)
}
version, err := sarama.ParseKafkaVersion(version)
if err != nil {
log.Panicf("Error parsing Kafka version: %v", err)
}
/**
* Construct a new Sarama configuration.
* The Kafka cluster version has to be defined before the consumer/producer is initialized.
*/
config := sarama.NewConfig()
config.Version = version
switch assignor {
case "sticky":
config.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.NewBalanceStrategySticky()}
case "roundrobin":
config.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.NewBalanceStrategyRoundRobin()}
case "range":
config.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.NewBalanceStrategyRange()}
default:
log.Panicf("Unrecognized consumer group partition assignor: %s", assignor)
}
if oldest {
config.Consumer.Offsets.Initial = sarama.OffsetOldest
}
/**
* Setup a new Sarama consumer group
*/
consumer := Consumer{
ready: make(chan bool),
}
ctx, cancel := context.WithCancel(context.Background())
client, err := sarama.NewConsumerGroup(strings.Split(brokers, ","), group, config)
if err != nil {
log.Panicf("Error creating consumer group client: %v", err)
}
consumptionIsPaused := false
wg := &sync.WaitGroup{}
wg.Add(1)
go func() {
defer wg.Done()
for {
// `Consume` should be called inside an infinite loop, when a
// server-side rebalance happens, the consumer session will need to be
// recreated to get the new claims
if err := client.Consume(ctx, strings.Split(topics, ","), &consumer); err != nil {
if errors.Is(err, sarama.ErrClosedConsumerGroup) {
return
}
log.Panicf("Error from consumer: %v", err)
}
// check if context was cancelled, signaling that the consumer should stop
if ctx.Err() != nil {
return
}
consumer.ready = make(chan bool)
}
}()
<-consumer.ready // Await till the consumer has been set up
log.Println("Sarama consumer up and running!...")
sigusr1 := make(chan os.Signal, 1)
signal.Notify(sigusr1, syscall.SIGUSR1)
sigterm := make(chan os.Signal, 1)
signal.Notify(sigterm, syscall.SIGINT, syscall.SIGTERM)
for keepRunning {
select {
case <-ctx.Done():
log.Println("terminating: context cancelled")
keepRunning = false
case <-sigterm:
log.Println("terminating: via signal")
keepRunning = false
case <-sigusr1:
toggleConsumptionFlow(client, &consumptionIsPaused)
}
}
cancel()
wg.Wait()
if err = client.Close(); err != nil {
log.Panicf("Error closing client: %v", err)
}
}
// / ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages().
func (consumer *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
var wg sync.WaitGroup
workerChan := make(chan *sarama.ConsumerMessage, 10)
// Start workers
for i := 0; i < 10; i++ {
wg.Add(1)
go func() {
defer wg.Done()
consumer.processMessages(workerChan, session)
}()
}
// Start ticker for reporting messages per second
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
messageCount := 0
go func() {
for range ticker.C {
log.Printf("Messages consumed per second: %d", messageCount)
messageCount = 0 // Reset count after reporting
}
}()
// Send messages to workerChan
for msg := range claim.Messages() {
workerChan <- msg
messageCount++
}
close(workerChan)
wg.Wait()
return nil
}
func (consumer *Consumer) processMessages(workerChan chan *sarama.ConsumerMessage, session sarama.ConsumerGroupSession) {
batch := make([]string, 0, batchSize)
timer := time.NewTimer(batchTime)
defer timer.Stop()
for {
select {
case msg, ok := <-workerChan:
if !ok {
// Print batch size if any remaining messages
//if len(batch) > 0 {
// consumer.printBatchSize(batch, "Final batch before shutdown")
//}
//return
}
batch = append(batch, string(msg.Value))
session.MarkMessage(msg, "")
//if len(batch) >= batchSize {
// consumer.printBatchSize(batch, "Batch processed due to reaching batch size")
// batch = batch[:0] // Reset batch
// timer.Reset(batchTime)
//}
case <-timer.C:
//if len(batch) > 0 {
// consumer.printBatchSize(batch, "Batch processed due to timer expiration")
// batch = batch[:0] // Reset batch
//}
timer.Reset(batchTime)
}
}
}
func (consumer *Consumer) printBatchSize(batch []string, reason string) {
log.Printf("%s: Batch size: %d", reason, len(batch))
}
// Setup is run at the beginning of a new session, before ConsumeClaim
func (consumer *Consumer) Setup(sarama.ConsumerGroupSession) error {
// Mark the consumer as ready
close(consumer.ready)
return nil
}
// Cleanup is run at the end of a session, once all ConsumeClaim goroutines have exited
func (consumer *Consumer) Cleanup(sarama.ConsumerGroupSession) error {
return nil
}
func toggleConsumptionFlow(client sarama.ConsumerGroup, isPaused *bool) {
if *isPaused {
client.ResumeAll()
log.Println("Resuming consumption")
} else {
client.PauseAll()
log.Println("Pausing consumption")
}
*isPaused = !*isPaused
}
The vector consumer topping out at 100-200k events is pretty slow and a bottleneck for our applications. The go variant was able to consume hundreds of thousands of messages before easily catching up to the producer script:
go run consume.go --brokers=localhost:29092 --group="mygroup" --topics=conn
2025/04/28 23:06:42 Starting a new Sarama consumer
2025/04/28 23:06:46 Sarama consumer up and running!...
2025/04/28 23:06:47 Messages consumed per second: 1033115
2025/04/28 23:06:48 Messages consumed per second: 1249705
2025/04/28 23:06:49 Messages consumed per second: 1164359
2025/04/28 23:06:50 Messages consumed per second: 1379555
2025/04/28 23:06:51 Messages consumed per second: 1300056
2025/04/28 23:06:52 Messages consumed per second: 1243945
2025/04/28 23:06:53 Messages consumed per second: 1248192
2025/04/28 23:06:54 Messages consumed per second: 1378353
2025/04/28 23:06:55 Messages consumed per second: 1119811
2025/04/28 23:06:56 Messages consumed per second: 1401732
2025/04/28 23:06:57 Messages consumed per second: 2611465
2025/04/28 23:06:58 Messages consumed per second: 1097094
2025/04/28 23:06:59 Messages consumed per second: 1380042
2025/04/28 23:07:00 Messages consumed per second: 1100183
2025/04/28 23:07:01 Messages consumed per second: 1026614
2025/04/28 23:07:02 Messages consumed per second: 1318862
2025/04/28 23:07:03 Messages consumed per second: 1147891
2025/04/28 23:07:04 Messages consumed per second: 1167338
2025/04/28 23:07:05 Messages consumed per second: 883569
2025/04/28 23:07:06 Messages consumed per second: 594106
2025/04/28 23:07:07 Messages consumed per second: 1343377
2025/04/28 23:07:08 Messages consumed per second: 1254939
2025/04/28 23:07:09 Messages consumed per second: 1257055
2025/04/28 23:07:10 Messages consumed per second: 1135398
2025/04/28 23:07:11 Messages consumed per second: 1043541
2025/04/28 23:07:12 Messages consumed per second: 878260
2025/04/28 23:07:13 Messages consumed per second: 1358307
2025/04/28 23:07:14 Messages consumed per second: 1351445
2025/04/28 23:07:15 Messages consumed per second: 1357376
2025/04/28 23:07:16 Messages consumed per second: 1083640
2025/04/28 23:07:17 Messages consumed per second: 1034354
2025/04/28 23:07:18 Messages consumed per second: 974207
2025/04/28 23:07:19 Messages consumed per second: 752598
2025/04/28 23:07:20 Messages consumed per second: 770773
2025/04/28 23:07:21 Messages consumed per second: 1293191
2025/04/28 23:07:22 Messages consumed per second: 1034238
2025/04/28 23:07:23 Messages consumed per second: 1204225
2025/04/28 23:07:24 Messages consumed per second: 1206141
2025/04/28 23:07:25 Messages consumed per second: 1033483
2025/04/28 23:07:26 Messages consumed per second: 1034252
2025/04/28 23:07:27 Messages consumed per second: 734936
2025/04/28 23:07:28 Messages consumed per second: 936752
2025/04/28 23:07:29 Messages consumed per second: 866167
2025/04/28 23:07:30 Messages consumed per second: 907015
2025/04/28 23:07:31 Messages consumed per second: 748981
2025/04/28 23:07:32 Messages consumed per second: 630329
2025/04/28 23:07:33 Messages consumed per second: 868474
2025/04/28 23:07:34 Messages consumed per second: 853221
2025/04/28 23:07:35 Messages consumed per second: 747441
2025/04/28 23:07:36 Messages consumed per second: 932502
2025/04/28 23:07:37 Messages consumed per second: 755004
2025/04/28 23:07:38 Messages consumed per second: 906781
2025/04/28 23:07:39 Messages consumed per second: 796320
2025/04/28 23:07:40 Messages consumed per second: 905314
2025/04/28 23:07:41 Messages consumed per second: 868931
2025/04/28 23:07:42 Messages consumed per second: 983289
2025/04/28 23:07:43 Messages consumed per second: 863403
2025/04/28 23:07:44 Messages consumed per second: 859279
2025/04/28 23:07:45 Messages consumed per second: 728868
2025/04/28 23:07:46 Messages consumed per second: 831269
2025/04/28 23:07:47 Messages consumed per second: 755107
2025/04/28 23:07:48 Messages consumed per second: 786263
2025/04/28 23:07:49 Messages consumed per second: 689819
2025/04/28 23:07:50 Messages consumed per second: 589625
2025/04/28 23:07:51 Messages consumed per second: 789969
2025/04/28 23:07:52 Messages consumed per second: 689350
2025/04/28 23:07:53 Messages consumed per second: 752703
2025/04/28 23:07:54 Messages consumed per second: 811179
2025/04/28 23:07:55 Messages consumed per second: 728849
2025/04/28 23:07:56 Messages consumed per second: 794851
2025/04/28 23:07:57 Messages consumed per second: 707884
2025/04/28 23:07:58 Messages consumed per second: 720693
2025/04/28 23:07:59 Messages consumed per second: 686657
2025/04/28 23:08:00 Messages consumed per second: 658642
2025/04/28 23:08:01 Messages consumed per second: 715750
2025/04/28 23:08:02 Messages consumed per second: 761442
2025/04/28 23:08:03 Messages consumed per second: 769389
2025/04/28 23:08:04 Messages consumed per second: 789766
2025/04/28 23:08:05 Messages consumed per second: 760283
2025/04/28 23:08:06 Messages consumed per second: 749480
2025/04/28 23:08:07 Messages consumed per second: 833138
2025/04/28 23:08:08 Messages consumed per second: 797727
2025/04/28 23:08:09 Messages consumed per second: 739036
2025/04/28 23:08:10 Messages consumed per second: 750551
2025/04/28 23:08:11 Messages consumed per second: 815162
2025/04/28 23:08:12 Messages consumed per second: 830347
2025/04/28 23:08:13 Messages consumed per second: 539575
2025/04/28 23:08:14 Messages consumed per second: 228995
2025/04/28 23:08:15 Messages consumed per second: 630482
2025/04/28 23:08:16 Messages consumed per second: 714526
2025/04/28 23:08:17 Messages consumed per second: 660833
2025/04/28 23:08:18 Messages consumed per second: 767530
2025/04/28 23:08:19 Messages consumed per second: 772910
2025/04/28 23:08:20 Messages consumed per second: 774460
2025/04/28 23:08:21 Messages consumed per second: 767329
2025/04/28 23:08:22 Messages consumed per second: 689275
2025/04/28 23:08:23 Messages consumed per second: 826536
2025/04/28 23:08:24 Messages consumed per second: 840705
2025/04/28 23:08:25 Messages consumed per second: 747046
2025/04/28 23:08:26 Messages consumed per second: 710926
2025/04/28 23:08:27 Messages consumed per second: 748926
2025/04/28 23:08:28 Messages consumed per second: 777634
2025/04/28 23:08:29 Messages consumed per second: 797520
2025/04/28 23:08:30 Messages consumed per second: 754631
2025/04/28 23:08:31 Messages consumed per second: 684700
2025/04/28 23:08:32 Messages consumed per second: 774795
2025/04/28 23:08:33 Messages consumed per second: 771609
2025/04/28 23:08:34 Messages consumed per second: 850415
2025/04/28 23:08:35 Messages consumed per second: 694146
2025/04/28 23:08:36 Messages consumed per second: 724833
2025/04/28 23:08:37 Messages consumed per second: 796513
2025/04/28 23:08:38 Messages consumed per second: 832136
2025/04/28 23:08:39 Messages consumed per second: 749271
2025/04/28 23:08:40 Messages consumed per second: 780125
2025/04/28 23:08:41 Messages consumed per second: 768679
2025/04/28 23:08:42 Messages consumed per second: 785076
2025/04/28 23:08:43 Messages consumed per second: 768122
2025/04/28 23:08:44 Messages consumed per second: 343500
2025/04/28 23:08:45 Messages consumed per second: 447239
2025/04/28 23:08:46 Messages consumed per second: 793903
2025/04/28 23:08:47 Messages consumed per second: 756941
2025/04/28 23:08:48 Messages consumed per second: 758989
2025/04/28 23:08:49 Messages consumed per second: 688653
2025/04/28 23:08:50 Messages consumed per second: 752550
2025/04/28 23:08:51 Messages consumed per second: 767974
2025/04/28 23:08:52 Messages consumed per second: 766501
2025/04/28 23:08:53 Messages consumed per second: 770079
2025/04/28 23:08:54 Messages consumed per second: 729111
2025/04/28 23:08:55 Messages consumed per second: 685869
2025/04/28 23:08:56 Messages consumed per second: 728583
2025/04/28 23:08:57 Messages consumed per second: 379380
2025/04/28 23:08:58 Messages consumed per second: 149988
2025/04/28 23:08:59 Messages consumed per second: 151153
2025/04/28 23:09:00 Messages consumed per second: 148993
2025/04/28 23:09:01 Messages consumed per second: 149775
2025/04/28 23:09:02 Messages consumed per second: 150662
2025/04/28 23:09:03 Messages consumed per second: 149320
2025/04/28 23:09:04 Messages consumed per second: 150468
2025/04/28 23:09:05 Messages consumed per second: 149729
2025/04/28 23:09:06 Messages consumed per second: 150349
2025/04/28 23:09:07 Messages consumed per second: 149430
2025/04/28 23:09:08 Messages consumed per second: 150104
2025/04/28 23:09:09 Messages consumed per second: 150816
2025/04/28 23:09:10 Messages consumed per second: 149357
2025/04/28 23:09:11 Messages consumed per second: 150757
2025/04/28 23:09:12 Messages consumed per second: 149164
2025/04/28 23:09:13 Messages consumed per second: 150540
2025/04/28 23:09:14 Messages consumed per second: 149378
2025/04/28 23:09:15 Messages consumed per second: 150678
2025/04/28 23:09:16 Messages consumed per second: 148783
2025/04/28 23:09:17 Messages consumed per second: 151275
Are there config options I should try? Is there a performance baseline for the kafka source available?
Configuration
[api]
enabled = true
[sources.kafka]
type = "kafka"
bootstrap_servers = "localhost:29092"
group_id = "aeohkh4k1j"
topics = ["conn"]
auto_offset_reset = "latest"
commit_interval_ms = 5000
metrics.topic_lag_metric = true
session_timeout_ms = 30000
[sinks.my_sink_id]
type = "blackhole"
inputs = [ "kafka" ]
Version
vector 0.46.1 (aarch64-apple-darwin 9a19e8a 2025-04-14 18:36:30.707862743)
Debug Output
Example Data
No response
Additional Context
Vector, Kafka running locally on a macbook pro, though I have replicated this on a 16cpu 32GB RAM EC2 instance as well
References
No response
A note for the community
Problem
Scenario:
A single node kafka broker, with 1 topic and one partition with 1 replica. Example docker-compose:
Vector source sending to the blackhole sink, example config:
I wrote a short python script to publish 150k messages per second. (just
pip install confluent-kafkaas a prereq):I observed through
vector topthat the kafka consume side could only consume 100-200k messages at best, even though hardware resources were available. No amount oflibrdkafka_optionstweaking could get better performance.To prove better performance was available on the same hardware I wrote the following go program:
The vector consumer topping out at 100-200k events is pretty slow and a bottleneck for our applications. The go variant was able to consume hundreds of thousands of messages before easily catching up to the producer script:
Are there config options I should try? Is there a performance baseline for the kafka source available?
Configuration
Version
vector 0.46.1 (aarch64-apple-darwin 9a19e8a 2025-04-14 18:36:30.707862743)
Debug Output
Example Data
No response
Additional Context
Vector, Kafka running locally on a macbook pro, though I have replicated this on a 16cpu 32GB RAM EC2 instance as well
References
No response