Kafka Python vs Confluent Kafka: Production Benchmarks
If you are building a new Python service that talks to Kafka today use confluent-kafka. It wraps librdkafka so you get C level performance thread safety and exactly once semantics without fighting the GIL. The pure Python kafka-python library is fine for low volume tooling but it falls over under load and lacks transactional producer support.
I learned this the hard way migrating an event ingestion pipeline that processes 50k messages per second. The old kafka-python consumers lagged during rebalances and the producer blocked the event loop during network hiccups. Switching to confluent-kafka cut p99 latency by 60 percent and eliminated the rebalance storms. Here is the data and the migration path so you do not repeat my mistakes.
Performance benchmarks and throughput comparison
The short answer: confluent-kafka sustains 3x to 5x higher throughput at lower CPU because the heavy lifting happens in librdkafka outside the GIL. kafka-python serializes and processes everything in the interpreter.
I ran a controlled test on a c6i.2xlarge producer sending 1KB JSON payloads to a three broker MSK cluster with acks=all.
| Client | Throughput (msg/s) | CPU User % | p99 Latency (ms) |
|---|---|---|---|
| kafka-python 2.0.2 | 42,000 | 85 | 145 |
| confluent-kafka 2.3.0 | 185,000 | 32 | 18 |
The confluent-kafka producer batches aggressively by default. The kafka-python producer requires manual tuning of batch_size linger_ms and buffer_memory to get close and even then it hits the GIL ceiling. For consumers the gap widens during rebalances because kafka-python stops the world in Python while confluent-kafka handles the protocol in C.
If you are running async FastAPI services the difference is stark. kafka-python has no native async support so you end up running consumers in threads or using aiokafka which adds another dependency layer. I covered the async FastAPI pattern in Kafka Python Tutorial: From Docker to FastAPI but that tutorial uses aiokafka. For new projects I would point you straight to confluent-kafka with a thread pool executor.
API design differences and developer experience
kafka-python feels Pythonic. The API mirrors the Java client closely with KafkaProducer KafkaConsumer and TopicPartition objects. Configuration is a dict. It is readable until you need advanced features.
confluent-kafka exposes the librdkafka surface area. Configuration uses dotted string keys like bootstrap.servers enable.auto.commit and transactional.id. It feels like C wrapped in Python because it is. The callback model for delivery reports and rebalance listeners is verbose but explicit.
Producer send comparison:
# kafka-pythonfrom kafka import KafkaProducerimport json
producer = KafkaProducer( bootstrap_servers="localhost:9092", value_serializer=lambda v: json.dumps(v).encode("utf-8"), acks="all", retries=3,)future = producer.send("events", {"user_id": 123, "action": "click"})record_metadata = future.get(timeout=10)print(f"partition {record_metadata.partition} offset {record_metadata.offset}")# confluent-kafkafrom confluent_kafka import Producerimport json
conf = { "bootstrap.servers": "localhost:9092", "acks": "all", "retries": 3, "linger.ms": 5,}producer = Producer(conf)
def delivery_report(err, msg): if err: print(f"Delivery failed: {err}") else: print(f"Delivered to {msg.topic()} [{msg.partition()}] @ {msg.offset()}")
producer.produce( "events", key=b"123", value=json.dumps({"user_id": 123, "action": "click"}).encode("utf-8"), callback=delivery_report,)producer.poll(0) # serve delivery callbacksproducer.flush(timeout=10)The confluent-kafka version requires a poll loop or a background thread to fire callbacks. In FastAPI I run poll in a lifespan startup task or a dedicated thread. The upside is backpressure visibility. The downside is more boilerplate.
Consumer side kafka-python gives you an iterator. confluent-kafka gives you poll(timeout) returning a message or None. You own the loop.
# confluent-kafka consumer with manual offset commitfrom confluent_kafka import Consumer, KafkaError
conf = { "bootstrap.servers": "localhost:9092", "group.id": "analytics-consumer", "enable.auto.commit": False, "auto.offset.reset": "earliest",}consumer = Consumer(conf)consumer.subscribe(["events"])
try: while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: continue raise KafkaException(msg.error()) process(msg.value()) consumer.commit(asynchronous=False)finally: consumer.close()This verbosity pays off when you need exactly once semantics.
Operational maturity: monitoring rebalancing and exactly-once semantics
kafka-python does not support idempotent producers or transactions. You can simulate idempotence with enable.idempotence=True in confluent-kafka and get exactly once semantics across partition boundaries using init_transactions begin_transaction send_offsets_to_transaction and commit_transaction.
# Exactly once producer consumer pattern confluent-kafkafrom confluent_kafka import Producer, Consumer, KafkaError, KafkaException
producer_conf = { "bootstrap.servers": "localhost:9092", "transactional.id": "etl-pipeline-1", "enable.idempotence": True, "acks": "all",}producer = Producer(producer_conf)producer.init_transactions()
consumer_conf = { "bootstrap.servers": "localhost:9092", "group.id": "etl-consumer", "enable.auto.commit": False, "isolation.level": "read_committed",}consumer = Consumer(consumer_conf)consumer.subscribe(["raw-events"])
try: while True: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): continue
producer.begin_transaction() try: transformed = transform(msg.value()) producer.produce("transformed-events", transformed) producer.send_offsets_to_transaction( [msg.topic_partition()], consumer.consumer_group_metadata() ) producer.commit_transaction() except Exception: producer.abort_transaction() raisefinally: consumer.close()This pattern is impossible with kafka-python. If you need exactly once you are blocked.
Rebalancing behavior differs too. kafka-python triggers a full stop the world rebalance in Python. confluent-kafka uses the cooperative sticky assignor (since librdkafka 1.6) which preserves partition ownership during rolling restarts. In production that means zero consumer lag spikes during deploys. I have seen kafka-python consumers take 30 seconds to rejoin a 20 partition topic. confluent-kafka rejoins in under 2 seconds.
Monitoring: confluent-kafka exposes stats_cb with a JSON blob containing broker level metrics queue sizes and throttle times. Pipe that to Datadog or Prometheus. kafka-python exposes nothing native. You have to instrument it yourself.
Migration guide: switching from kafka-python to confluent-kafka
Step 1: Add confluent-kafka to requirements. Keep kafka-python installed during the cutover.
Step 2: Rewrite producers first. They are stateless. Use the delivery callback pattern. Wrap the producer in a class that matches your existing interface.
# wrapper to match existing kafka-python interfaceclass KafkaProducerWrapper: def __init__(self, bootstrap_servers: str): self._producer = Producer({ "bootstrap.servers": bootstrap_servers, "acks": "all", "enable.idempotence": True, "linger.ms": 5, })
def send(self, topic: str, value: dict, key: str = None): from concurrent.futures import Future future = Future()
def callback(err, msg): if err: future.set_exception(KafkaException(err)) else: future.set_result(msg)
self._producer.produce( topic, key=key.encode() if key else None, value=json.dumps(value).encode(), callback=callback, ) self._producer.poll(0) return future
def flush(self, timeout: float = 10): self._producer.flush(timeout)
def close(self): self._producer.flush()Step 3: Rewrite consumers. This is harder because of offset management. Run both consumers in parallel for a week writing to a shadow topic. Compare offsets and latency. Cut traffic only after parity.
Step 4: Update your health checks. confluent-kafka does not have a bootstrap_connected method. Use list_topics(timeout=5) as a liveness probe.
Step 5: Remove kafka-python and aiokafka from requirements.
I followed this path on a payments ingestion service. The shadow run caught a serialization bug where None keys behaved differently. kafka-python sends null key. confluent-kafka requires explicit key=None or omission. The wrapper handles it.
Decision matrix: when to choose which client for production
Choose confluent-kafka when:
- Throughput exceeds 10k msg/s sustained
- You need exactly once semantics
- You run in Kubernetes with rolling restarts
- You need cooperative rebalancing
- You want native metrics without sidecars
- You are building new services
Choose kafka-python when:
- You maintain a legacy codebase with no migration budget
- Volume is under 5k msg/s and latency is not critical
- You need pure Python for WASM or constrained environments
- Your team refuses C extensions (rare but happens)
Choose aiokafka when:
- You are deep in asyncio and cannot introduce threads
- You accept the performance ceiling and maintenance risk
I default to confluent-kafka for every new service. The operational savings on rebalances alone justify the learning curve. If you are building async FastAPI APIs with SQLAlchemy 2.0 the pattern fits cleanly with a thread pool executor for the producer poll loop. See Fix AsyncSessionLocal Errors in FastAPI (SQLAlchemy 2.0) for the async DB layer that pairs well with this.
FAQ
Is confluent-kafka thread safe?
Yes. The Producer is thread safe for produce and flush. The Consumer is not thread safe. Use one consumer per thread or process.
Does confluent-kafka work with asyncio?
Not natively. Run poll in a background thread or use run_in_executor. The wrapper pattern above returns a Future compatible with asyncio.
What about schema registry integration?
Both clients work with confluent-kafka-schema-registry or fastavro. confluent-kafka has first party Avro and Protobuf serializers in confluent_kafka.schema_registry.
Can I use kafka-python and confluent-kafka in the same process? Yes. They use different underlying libraries. No symbol conflicts. Useful during migration.
Key Takeaways
- confluent-kafka wraps librdkafka delivering 3x to 5x throughput at lower CPU
- kafka-python lacks exactly once transactions and cooperative rebalancing
- Migration is low risk if you wrap producers and shadow run consumers
- Default to confluent-kafka for new production services
- Use a thread pool executor for the poll loop in async FastAPI applications
Working on something similar?
If you're building backend or AI systems and want a second set of senior eyes, let's talk.