Apache Kafka: Producers & Consumers
Producer
# pip install confluent-kafka
from confluent_kafka import Producer
import json
producer = Producer({
'bootstrap.servers': 'localhost:9092',
'acks': 'all', # wait for all ISR replicas
'retries': 5,
'retry.backoff.ms': 300,
'linger.ms': 5, # batch messages for 5ms (higher throughput)
'batch.size': 16384,
'compression.type': 'snappy', # compress batches (lz4, gzip, snappy, zstd)
'enable.idempotence': True, # exactly-once delivery
})
def delivery_callback(err, msg):
if err:
print(f'Delivery failed: {err}')
else:
print(f'Delivered to {msg.topic()}[{msg.partition()}] at offset {msg.offset()}')
# Produce with key (same key → same partition = ordered per key)
order = {'orderId': 1, 'userId': 'user-123', 'amount': 99.99}
producer.produce(
topic='orders',
key='user-123', # key ensures ordering per user
value=json.dumps(order).encode('utf-8'),
callback=delivery_callback,
)
producer.poll(0) # trigger delivery callbacks (non-blocking)
# Flush before exit
producer.flush(timeout=10) # wait for all pending messagesConsumer
from confluent_kafka import Consumer, KafkaError
import json
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'order-processing-group',
'auto.offset.reset': 'earliest', # 'latest' = only new messages
'enable.auto.commit': False, # manual commit for at-least-once
'max.poll.interval.ms': 300000, # max time between polls before rebalance
'session.timeout.ms': 45000,
})
consumer.subscribe(['orders'])
try:
while True:
msg = consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
continue # reached end of partition
else:
raise KafkaException(msg.error())
# Process message
key = msg.key().decode('utf-8') if msg.key() else None
value = json.loads(msg.value().decode('utf-8'))
print(f'Offset {msg.offset()}: key={key}, value={value}')
process_order(value)
# Commit after successful processing (at-least-once delivery)
consumer.commit(asynchronous=False)
except KeyboardInterrupt:
pass
finally:
consumer.close() # commits offsets and triggers rebalanceJava Producer/Consumer
// Maven: kafka-clients dependency
// Producer
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
ProducerRecord<String, String> record = new ProducerRecord<>(
"orders", "user-123", "{"orderId": 1}"
);
producer.send(record, (metadata, exception) -> {
if (exception != null) log.error("Send failed", exception);
else log.info("Sent to partition {} offset {}", metadata.partition(), metadata.offset());
});
}
// Consumer
Properties cProps = new Properties();
cProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
cProps.put(ConsumerConfig.GROUP_ID_CONFIG, "order-group");
cProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
cProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
cProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
cProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(cProps)) {
consumer.subscribe(List.of("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
processOrder(record.key(), record.value());
}
consumer.commitSync();
}
}Keep your own version of these notes — editable, searchable, and organised by your stack.
Start free