Skip to content

Poor kafka source consumer performance #22958

Description

@john-from-corelight

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.

Image

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    domain: performanceAnything related to Vector's performancesource: kafkaAnything `kafka` source related

    Type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions