Apache Kafka
02 / 03

Producers & Consumers

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 messages

Consumer

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 rebalance

Java 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