Partitions Enable Parallelism
Partitions allow multiple consumers to process topic in parallel. Ordering guaranteed per partition.
Apache Kafka is a distributed streaming platform designed for building real-time data pipelines and streaming applications. Originally developed by LinkedIn and open-sourced in 2011, Kafka has become the de facto standard for event streaming in modern distributed systems.
Kafka operates as a distributed commit log - a persistent, append-only data structure that stores streams of records (events) in topics. Unlike traditional message queues that delete messages after consumption, Kafka retains messages for a configurable retention period, allowing multiple consumers to read the same messages at different times and speeds.
This design makes Kafka ideal for scenarios where you need to:
A topic is a category or stream of messages in Kafka. Think of it as a named channel where producers publish messages and consumers subscribe to read them. Topics are similar to tables in a database or folders in a filesystem - they provide logical organization for related messages.
Key Characteristics:
Topics in Kafka are immutable logs - messages are appended to the end and never modified or deleted (until retention expires). This append-only design provides several benefits: it enables efficient sequential disk I/O, allows multiple consumers to read the same messages independently, and supports event replay for debugging or reprocessing.
Each topic is partitioned for scalability and replicated across multiple brokers for fault tolerance. Messages within a partition are strictly ordered, but ordering across partitions is not guaranteed unless you use a partitioning key.
Characteristics:
A partition is an ordered, immutable sequence of messages within a topic. Each topic is divided into one or more partitions, which enables Kafka to scale horizontally and process messages in parallel.
Why Partitions Matter:
Partitions are Kafka’s fundamental unit of parallelism. Without partitions, a topic would be processed by only one consumer at a time, creating a bottleneck. With partitions, different consumers can process different partitions simultaneously, dramatically increasing throughput.
Partitioning Strategy:
When a producer sends a message, Kafka determines which partition to use based on:
Real-World Example: In an e-commerce system, you might partition the orders topic by customer_id. This ensures all orders from the same customer are processed in order, while orders from different customers can be processed in parallel.
Why Partitions?
Partitioning Strategy:
A consumer group is a set of consumers that work together to consume messages from one or more topics. Kafka automatically distributes partitions across consumers in the same group, ensuring each partition is consumed by exactly one consumer in the group.
How Consumer Groups Work:
When consumers in a group subscribe to a topic, Kafka performs partition assignment - it divides the topic’s partitions among the available consumers. If you have 3 partitions and 2 consumers, Kafka might assign partitions 0 and 1 to consumer 1, and partition 2 to consumer 2.
Rebalancing:
When consumers join or leave a consumer group, Kafka automatically rebalances - it redistributes partitions among the remaining consumers. During rebalancing, consumers stop processing messages, which can cause a brief pause. This is why it’s important to design consumers to handle rebalancing gracefully and commit offsets frequently.
Scaling Pattern:
Key Rules:
Example:
An offset is a sequential number that uniquely identifies the position of a message within a partition. Offsets start at 0 and increment for each message. They are immutable - once assigned, an offset never changes.
Offset Management:
Consumers track their current offset - the position of the last message they’ve successfully processed. After processing a message, the consumer commits the offset, which tells Kafka “I’ve processed up to this point.” If the consumer crashes and restarts, it resumes from the last committed offset, ensuring no messages are lost and no messages are processed twice (assuming proper offset management).
Offset Commit Strategies:
Offset Storage:
Kafka stores committed offsets in a special internal topic called __consumer_offsets. This allows Kafka to track consumer progress even if consumers restart or rebalance.
Offset Management:
Broker = Kafka server. Cluster = Multiple brokers.
Replication:
from kafka import KafkaProducerimport json
class KafkaEventProducer: """Kafka producer for events"""
def __init__(self, bootstrap_servers: list): self.producer = KafkaProducer( bootstrap_servers=bootstrap_servers, value_serializer=lambda v: json.dumps(v).encode('utf-8'), key_serializer=lambda k: k.encode('utf-8') if k else None, # Idempotent producer (exactly-once) enable_idempotence=True, acks='all', # Wait for all replicas retries=3, max_in_flight_requests_per_connection=1 )
def send_event(self, topic: str, event: dict, key: str = None): """Send event to topic""" future = self.producer.send( topic, key=key, value=event )
# Wait for acknowledgment try: record_metadata = future.get(timeout=10) print(f"Event sent to {record_metadata.topic} " f"partition {record_metadata.partition} " f"offset {record_metadata.offset}") return record_metadata except Exception as e: print(f"Error sending event: {e}") raise
def send_with_callback(self, topic: str, event: dict, key: str = None): """Send event with callback""" def on_send_success(record_metadata): print(f"Event sent: {record_metadata.topic}/" f"{record_metadata.partition}/" f"{record_metadata.offset}")
def on_send_error(exception): print(f"Error sending event: {exception}")
self.producer.send( topic, key=key, value=event ).add_callback(on_send_success).add_errback(on_send_error)
def flush(self): """Flush pending messages""" self.producer.flush()
def close(self): """Close producer""" self.producer.close()
# Usageproducer = KafkaEventProducer(['localhost:9092'])
# Send eventproducer.send_event('user-events', { 'event_type': 'user.created', 'user_id': 123, 'timestamp': '2024-01-01T10:00:00Z'}, key='123') # Key ensures same user goes to same partition
producer.flush()producer.close()import org.apache.kafka.clients.producer.*;import org.apache.kafka.common.serialization.StringSerializer;import com.fasterxml.jackson.databind.ObjectMapper;import java.util.Properties;
public class KafkaEventProducer { private final KafkaProducer<String, String> producer; private final ObjectMapper objectMapper = new ObjectMapper();
public KafkaEventProducer(String bootstrapServers) { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
// Idempotent producer (exactly-once) props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 1);
this.producer = new KafkaProducer<>(props); }
public void sendEvent(String topic, Map<String, Object> event, String key) { try { String value = objectMapper.writeValueAsString(event);
ProducerRecord<String, String> record = new ProducerRecord<>( topic, key, value );
// Send with callback producer.send(record, (metadata, exception) -> { if (exception != null) { System.err.println("Error sending event: " + exception.getMessage()); } else { System.out.println("Event sent to " + metadata.topic() + "/" + metadata.partition() + "/" + metadata.offset()); } }); } catch (Exception e) { System.err.println("Error serializing event: " + e.getMessage()); } }
public void flush() { producer.flush(); }
public void close() { producer.close(); }}
// UsageKafkaEventProducer producer = new KafkaEventProducer("localhost:9092");
Map<String, Object> event = new HashMap<>();event.put("event_type", "user.created");event.put("user_id", 123);event.put("timestamp", "2024-01-01T10:00:00Z");
producer.sendEvent("user-events", event, "123");producer.flush();producer.close();from kafka import KafkaConsumerfrom kafka.errors import KafkaErrorimport jsonfrom typing import Callable
class KafkaEventConsumer: """Kafka consumer for events"""
def __init__(self, bootstrap_servers: list, group_id: str): self.consumer = KafkaConsumer( bootstrap_servers=bootstrap_servers, group_id=group_id, value_deserializer=lambda m: json.loads(m.decode('utf-8')), key_deserializer=lambda k: k.decode('utf-8') if k else None, enable_auto_commit=False, # Manual offset commit auto_offset_reset='earliest', # Start from beginning if no offset max_poll_records=100 # Batch size )
def subscribe(self, topics: list): """Subscribe to topics""" self.consumer.subscribe(topics) print(f"Subscribed to topics: {topics}")
def consume(self, handler: Callable): """Consume messages""" try: while True: # Poll for messages (batch) message_batch = self.consumer.poll(timeout_ms=1000)
for topic_partition, messages in message_batch.items(): for message in messages: try: # Process message handler(message.value, message.key, message.offset)
# Commit offset after processing self.consumer.commit() except Exception as e: print(f"Error processing message: {e}") # Don't commit - will retry except KeyboardInterrupt: print("Stopping consumer...") finally: self.consumer.close()
def consume_with_manual_commit(self, handler: Callable): """Consume with manual offset commit""" try: while True: message_batch = self.consumer.poll(timeout_ms=1000)
offsets_to_commit = {}
for topic_partition, messages in message_batch.items(): for message in messages: try: # Process message handler(message.value, message.key, message.offset)
# Track offset for commit offsets_to_commit[topic_partition] = \ OffsetAndMetadata(message.offset + 1, None) except Exception as e: print(f"Error processing message: {e}") # Don't commit failed messages break
# Commit all processed offsets if offsets_to_commit: self.consumer.commit(offsets_to_commit) except KeyboardInterrupt: print("Stopping consumer...") finally: self.consumer.close()
# Usagedef handle_event(event, key, offset): print(f"Processing event: {event['event_type']} " f"key: {key} offset: {offset}") # Process event...
consumer = KafkaEventConsumer( ['localhost:9092'], group_id='event-processors')consumer.subscribe(['user-events', 'order-events'])consumer.consume(handle_event)import org.apache.kafka.clients.consumer.*;import org.apache.kafka.common.serialization.StringDeserializer;import com.fasterxml.jackson.databind.ObjectMapper;import java.time.Duration;import java.util.*;
public class KafkaEventConsumer { private final KafkaConsumer<String, String> consumer; private final ObjectMapper objectMapper = new ObjectMapper();
public KafkaEventConsumer(String bootstrapServers, String groupId) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // Manual commit props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
this.consumer = new KafkaConsumer<>(props); }
public void subscribe(List<String> topics) { consumer.subscribe(topics); System.out.println("Subscribed to topics: " + topics); }
public void consume(java.util.function.Consumer<Map<String, Object>> handler) { try { while (true) { // Poll for messages (batch) ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) { try { // Deserialize event Map<String, Object> event = objectMapper.readValue( record.value(), Map.class );
// Process event handler.accept(event);
// Commit offset after processing consumer.commitSync(); } catch (Exception e) { System.err.println("Error processing message: " + e.getMessage()); // Don't commit - will retry } } } } catch (Exception e) { System.err.println("Consumer error: " + e.getMessage()); } finally { consumer.close(); } }}
// UsageKafkaEventConsumer consumer = new KafkaEventConsumer( "localhost:9092", "event-processors");
consumer.subscribe(Arrays.asList("user-events", "order-events"));
consumer.consume(event -> { System.out.println("Processing event: " + event.get("event_type")); // Process event...});import { Kafka, Consumer, EachMessagePayload } from 'kafkajs';
interface Event { event_type: string; [key: string]: any;}
class KafkaEventConsumer { private consumer: Consumer;
constructor(bootstrapServers: string[], groupId: string) { const kafka = new Kafka({ clientId: 'event-consumer', brokers: bootstrapServers, });
this.consumer = kafka.consumer({ groupId: groupId, maxBytesPerPartition: 1048576, // 1MB }); }
async connect(): Promise<void> { await this.consumer.connect(); }
async subscribe(topics: string[]): Promise<void> { // Subscribe to topics await this.consumer.subscribe({ topics }); console.log(`Subscribed to topics: ${topics.join(', ')}`); }
async consume(handler: (event: Event, key: string | null, offset: string) => Promise<void>): Promise<void> { // Consume messages await this.consumer.run({ eachMessage: async (payload: EachMessagePayload) => { try { const event: Event = JSON.parse(payload.message.value?.toString() || '{}'); const key = payload.message.key?.toString() || null; const offset = payload.message.offset;
// Process message await handler(event, key, offset);
// Offset is automatically committed } catch (error) { console.error(`Error processing message: ${error}`); // Don't commit - will retry } }, }); }
async close(): Promise<void> { // Close consumer await this.consumer.disconnect(); }}
// Usageasync function handleEvent(event: Event, key: string | null, offset: string) { console.log(`Processing event: ${event.event_type} key: ${key} offset: ${offset}`); // Process event...}
const consumer = new KafkaEventConsumer(['localhost:9092'], 'event-processors');await consumer.connect();await consumer.subscribe(['user-events', 'order-events']);await consumer.consume(handleEvent);#include <librdkafka/rdkafka.h>#include <iostream>#include <string>#include <nlohmann/json.hpp>
class KafkaEventConsumer {private: rd_kafka_t* consumer; rd_kafka_conf_t* conf;
public: KafkaEventConsumer(const std::vector<std::string>& bootstrapServers, const std::string& groupId) { char errstr[512];
// Create configuration conf = rd_kafka_conf_new();
// Set bootstrap servers std::string brokers; for (const auto& server : bootstrapServers) { if (!brokers.empty()) brokers += ","; brokers += server; } rd_kafka_conf_set(conf, "bootstrap.servers", brokers.c_str(), errstr, sizeof(errstr)); rd_kafka_conf_set(conf, "group.id", groupId.c_str(), errstr, sizeof(errstr)); rd_kafka_conf_set(conf, "enable.auto.commit", "false", errstr, sizeof(errstr)); rd_kafka_conf_set(conf, "auto.offset.reset", "earliest", errstr, sizeof(errstr));
// Create consumer consumer = rd_kafka_new(RD_KAFKA_CONSUMER, conf, errstr, sizeof(errstr)); if (!consumer) { throw std::runtime_error("Failed to create consumer: " + std::string(errstr)); } }
void subscribe(const std::vector<std::string>& topics) { // Subscribe to topics rd_kafka_topic_partition_list_t* topicList = rd_kafka_topic_partition_list_new(topics.size()); for (const auto& topic : topics) { rd_kafka_topic_partition_list_add(topicList, topic.c_str(), RD_KAFKA_PARTITION_UA); }
rd_kafka_resp_err_t err = rd_kafka_subscribe(consumer, topicList); if (err) { throw std::runtime_error("Failed to subscribe: " + std::string(rd_kafka_err2str(err))); }
rd_kafka_topic_partition_list_destroy(topicList); }
void consume(std::function<void(const nlohmann::json&, const std::string&, int64_t)> handler) { // Consume messages while (true) { rd_kafka_message_t* msg = rd_kafka_consumer_poll(consumer, 1000);
if (msg) { if (msg->err == RD_KAFKA_RESP_ERR_NO_ERROR) { try { nlohmann::json event = nlohmann::json::parse(static_cast<const char*>(msg->payload)); std::string key = msg->key ? std::string(static_cast<const char*>(msg->key), msg->key_len) : "";
// Process message handler(event, key, msg->offset);
// Commit offset after processing rd_kafka_commit_message(consumer, msg, 0); } catch (const std::exception& e) { std::cerr << "Error processing message: " << e.what() << std::endl; // Don't commit - will retry } }
rd_kafka_message_destroy(msg); } } }
~KafkaEventConsumer() { rd_kafka_consumer_close(consumer); rd_kafka_destroy(consumer); }};
// Usagevoid handleEvent(const nlohmann::json& event, const std::string& key, int64_t offset) { std::cout << "Processing event: " << event["event_type"] << " key: " << key << " offset: " << offset << std::endl; // Process event...}
int main() { KafkaEventConsumer consumer({"localhost:9092"}, "event-processors"); consumer.subscribe({"user-events", "order-events"}); consumer.consume(handleEvent); return 0;}using Confluent.Kafka;using System;using System.Collections.Generic;using System.Text.Json;using System.Threading;using System.Threading.Tasks;
public class KafkaEventConsumer { private readonly IConsumer<string, string> consumer;
public KafkaEventConsumer(string bootstrapServers, string groupId) { var config = new ConsumerConfig { BootstrapServers = bootstrapServers, GroupId = groupId, EnableAutoCommit = false, // Manual commit AutoOffsetReset = AutoOffsetReset.Earliest, MaxPartitionFetchBytes = 1048576 // 1MB };
consumer = new ConsumerBuilder<string, string>(config).Build(); }
public void Subscribe(IEnumerable<string> topics) { // Subscribe to topics consumer.Subscribe(topics); Console.WriteLine($"Subscribed to topics: {string.Join(", ", topics)}"); }
public void Consume(Func<object, string, long, Task> handler, CancellationToken cancellationToken) { // Consume messages try { while (!cancellationToken.IsCancellationRequested) { var result = consumer.Consume(cancellationToken);
try { // Deserialize event var eventData = JsonSerializer.Deserialize<object>(result.Message.Value);
// Process message handler(eventData, result.Message.Key, result.Offset).Wait();
// Commit offset after processing consumer.Commit(result); } catch (Exception e) { Console.Error.WriteLine($"Error processing message: {e.Message}"); // Don't commit - will retry } } } catch (OperationCanceledException) { Console.WriteLine("Stopping consumer..."); } finally { consumer.Close(); } }}
// Usageasync Task HandleEvent(object eventData, string key, long offset) { Console.WriteLine($"Processing event: {eventData} key: {key} offset: {offset}"); // Process event...}
var consumer = new KafkaEventConsumer("localhost:9092", "event-processors");consumer.Subscribe(new[] { "user-events", "order-events" });
var cts = new CancellationTokenSource();consumer.Consume(HandleEvent, cts.Token);Exactly-once semantics ensures that each message is processed exactly once, with no duplicates and no lost messages. This is the strongest delivery guarantee but also the most complex to implement.
In distributed systems, achieving exactly-once processing is challenging because:
Achieving exactly-once semantics requires coordination across three components:
Idempotent Producer
Transactional Producer
Idempotent Consumer
from kafka import KafkaProducer, KafkaConsumerfrom kafka.errors import KafkaErrorimport jsonfrom typing import Set
class ExactlyOnceProcessor: """Exactly-once message processing"""
def __init__(self, bootstrap_servers: list, group_id: str): # Idempotent producer self.producer = KafkaProducer( bootstrap_servers=bootstrap_servers, enable_idempotence=True, acks='all', transactional_id='exactly-once-producer' )
# Consumer with manual commit self.consumer = KafkaConsumer( bootstrap_servers=bootstrap_servers, group_id=group_id, enable_auto_commit=False, isolation_level='read_committed' # Only read committed messages )
# Track processed offsets (for idempotency) self.processed_offsets: Set[tuple] = set()
def process_exactly_once(self, topic: str, handler): """Process messages exactly once""" self.consumer.subscribe([topic])
try: while True: message_batch = self.consumer.poll(timeout_ms=1000)
for topic_partition, messages in message_batch.items(): # Begin transaction self.producer.begin_transaction()
try: for message in messages: # Check if already processed (idempotency) offset_key = (topic_partition.topic, topic_partition.partition, message.offset)
if offset_key in self.processed_offsets: print(f"Skipping duplicate: {offset_key}") continue
# Process message result = handler(message.value)
# Send result to output topic (in transaction) if result: self.producer.send('output-topic', value=json.dumps(result))
# Mark as processed self.processed_offsets.add(offset_key)
# Commit transaction (atomic) self.producer.commit_transaction()
# Commit consumer offset self.consumer.commit() except Exception as e: # Abort transaction on error self.producer.abort_transaction() print(f"Error processing: {e}") finally: self.producer.close() self.consumer.close()import org.apache.kafka.clients.producer.KafkaProducer;import org.apache.kafka.clients.producer.ProducerConfig;import org.apache.kafka.clients.consumer.KafkaConsumer;import org.apache.kafka.clients.consumer.ConsumerConfig;import org.apache.kafka.clients.consumer.ConsumerRecord;import org.apache.kafka.clients.consumer.ConsumerRecords;import java.util.*;
public class ExactlyOnceProcessor { private final KafkaProducer<String, String> producer; private final KafkaConsumer<String, String> consumer; private final Set<String> processedOffsets = new HashSet<>();
public ExactlyOnceProcessor(String bootstrapServers, String groupId) { // Idempotent producer Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "exactly-once-producer"); this.producer = new KafkaProducer<>(producerProps); this.producer.initTransactions();
// Consumer with manual commit Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); this.consumer = new KafkaConsumer<>(consumerProps); }
public void processExactlyOnce(String topic, java.util.function.Function<String, String> handler) { consumer.subscribe(Collections.singletonList(topic));
try { while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
// Begin transaction producer.beginTransaction();
try { for (ConsumerRecord<String, String> record : records) { // Check if already processed String offsetKey = record.topic() + "-" + record.partition() + "-" + record.offset();
if (processedOffsets.contains(offsetKey)) { System.out.println("Skipping duplicate: " + offsetKey); continue; }
// Process message String result = handler.apply(record.value());
// Send result (in transaction) if (result != null) { producer.send(new ProducerRecord<>("output-topic", result)); }
// Mark as processed processedOffsets.add(offsetKey); }
// Commit transaction producer.commitTransaction();
// Commit consumer offset consumer.commitSync(); } catch (Exception e) { // Abort transaction on error producer.abortTransaction(); System.err.println("Error processing: " + e.getMessage()); } } } finally { producer.close(); consumer.close(); } }}import { Kafka, Producer, Consumer, EachMessagePayload } from 'kafkajs';
interface Event { [key: string]: any;}
class ExactlyOnceProcessor { private producer: Producer; private consumer: Consumer; private processedOffsets: Set<string> = new Set();
constructor(bootstrapServers: string[], groupId: string) { const kafka = new Kafka({ clientId: 'exactly-once-processor', brokers: bootstrapServers, });
// Idempotent producer this.producer = kafka.producer({ idempotent: true, transactionalId: 'exactly-once-producer', maxInFlightRequests: 1, acks: -1, });
// Consumer with manual commit this.consumer = kafka.consumer({ groupId: groupId, readUncommitted: false, // Only read committed messages }); }
async connect(): Promise<void> { await this.producer.connect(); await this.consumer.connect(); await this.producer.transaction(); }
async processExactlyOnce( topic: string, handler: (event: Event) => Promise<Event | null> ): Promise<void> { await this.consumer.subscribe({ topics: [topic] });
await this.consumer.run({ eachMessage: async (payload: EachMessagePayload) => { // Check if already processed const offsetKey = `${payload.topic}-${payload.partition}-${payload.message.offset}`;
if (this.processedOffsets.has(offsetKey)) { console.log(`Skipping duplicate: ${offsetKey}`); return; }
try { // Begin transaction await this.producer.beginTransaction();
const event: Event = JSON.parse(payload.message.value?.toString() || '{}');
// Process message const result = await handler(event);
// Send result to output topic (in transaction) if (result) { await this.producer.send({ topic: 'output-topic', messages: [{ value: JSON.stringify(result), }], }); }
// Mark as processed this.processedOffsets.add(offsetKey);
// Commit transaction await this.producer.commitTransaction(); } catch (error) { // Abort transaction on error await this.producer.abortTransaction(); console.error(`Error processing: ${error}`); } }, }); }
async close(): Promise<void> { await this.producer.disconnect(); await this.consumer.disconnect(); }}
// Usageconst processor = new ExactlyOnceProcessor(['localhost:9092'], 'event-processors');await processor.connect();
await processor.processExactlyOnce('input-topic', async (event) => { console.log('Processing event:', event); // Process event... return { processed: true, event };});#include <librdkafka/rdkafka.h>#include <iostream>#include <unordered_set>#include <string>#include <nlohmann/json.hpp>
class ExactlyOnceProcessor {private: rd_kafka_t* producer; rd_kafka_t* consumer; std::unordered_set<std::string> processedOffsets;
public: ExactlyOnceProcessor(const std::vector<std::string>& bootstrapServers, const std::string& groupId) { char errstr[512];
// Idempotent producer configuration rd_kafka_conf_t* producerConf = rd_kafka_conf_new(); std::string brokers; for (const auto& server : bootstrapServers) { if (!brokers.empty()) brokers += ","; brokers += server; } rd_kafka_conf_set(producerConf, "bootstrap.servers", brokers.c_str(), errstr, sizeof(errstr)); rd_kafka_conf_set(producerConf, "enable.idempotence", "true", errstr, sizeof(errstr)); rd_kafka_conf_set(producerConf, "transactional.id", "exactly-once-producer", errstr, sizeof(errstr)); rd_kafka_conf_set(producerConf, "acks", "all", errstr, sizeof(errstr));
producer = rd_kafka_new(RD_KAFKA_PRODUCER, producerConf, errstr, sizeof(errstr)); rd_kafka_init_transactions(producer, -1);
// Consumer configuration rd_kafka_conf_t* consumerConf = rd_kafka_conf_new(); rd_kafka_conf_set(consumerConf, "bootstrap.servers", brokers.c_str(), errstr, sizeof(errstr)); rd_kafka_conf_set(consumerConf, "group.id", groupId.c_str(), errstr, sizeof(errstr)); rd_kafka_conf_set(consumerConf, "enable.auto.commit", "false", errstr, sizeof(errstr)); rd_kafka_conf_set(consumerConf, "isolation.level", "read_committed", errstr, sizeof(errstr));
consumer = rd_kafka_new(RD_KAFKA_CONSUMER, consumerConf, errstr, sizeof(errstr)); }
void processExactlyOnce(const std::string& topic, std::function<nlohmann::json(const nlohmann::json&)> handler) { // Subscribe to topic rd_kafka_topic_partition_list_t* topicList = rd_kafka_topic_partition_list_new(1); rd_kafka_topic_partition_list_add(topicList, topic.c_str(), RD_KAFKA_PARTITION_UA); rd_kafka_subscribe(consumer, topicList); rd_kafka_topic_partition_list_destroy(topicList);
while (true) { rd_kafka_message_t* msg = rd_kafka_consumer_poll(consumer, 1000);
if (msg && msg->err == RD_KAFKA_RESP_ERR_NO_ERROR) { // Check if already processed std::string offsetKey = std::string(msg->rkt->topic) + "-" + std::to_string(msg->partition) + "-" + std::to_string(msg->offset);
if (processedOffsets.find(offsetKey) != processedOffsets.end()) { std::cout << "Skipping duplicate: " << offsetKey << std::endl; rd_kafka_message_destroy(msg); continue; }
try { // Begin transaction rd_kafka_resp_err_t err = rd_kafka_begin_transaction(producer); if (err) { throw std::runtime_error("Failed to begin transaction"); }
nlohmann::json event = nlohmann::json::parse(static_cast<const char*>(msg->payload));
// Process message nlohmann::json result = handler(event);
// Send result (in transaction) if (!result.is_null()) { std::string resultStr = result.dump(); rd_kafka_producev( producer, RD_KAFKA_V_TOPIC("output-topic"), RD_KAFKA_V_VALUE(resultStr.c_str(), resultStr.length()), RD_KAFKA_V_END ); }
// Mark as processed processedOffsets.insert(offsetKey);
// Commit transaction rd_kafka_commit_transaction(producer, -1); } catch (const std::exception& e) { // Abort transaction on error rd_kafka_abort_transaction(producer, -1); std::cerr << "Error processing: " << e.what() << std::endl; }
rd_kafka_message_destroy(msg); } } }
~ExactlyOnceProcessor() { rd_kafka_destroy(producer); rd_kafka_consumer_close(consumer); rd_kafka_destroy(consumer); }};
// Usagenlohmann::json handleEvent(const nlohmann::json& event) { std::cout << "Processing event: " << event.dump() << std::endl; // Process event... return {{"processed", true}, {"event", event}};}
int main() { ExactlyOnceProcessor processor({"localhost:9092"}, "event-processors"); processor.processExactlyOnce("input-topic", handleEvent); return 0;}using Confluent.Kafka;using System;using System.Collections.Generic;using System.Text.Json;using System.Threading;using System.Threading.Tasks;
public class ExactlyOnceProcessor { private readonly IProducer<string, string> producer; private readonly IConsumer<string, string> consumer; private readonly HashSet<string> processedOffsets;
public ExactlyOnceProcessor(string bootstrapServers, string groupId) { // Idempotent producer var producerConfig = new ProducerConfig { BootstrapServers = bootstrapServers, EnableIdempotence = true, TransactionalId = "exactly-once-producer", Acks = Acks.All, MaxInFlightRequestsPerConnection = 1 }; producer = new ProducerBuilder<string, string>(producerConfig).Build(); producer.InitTransactions(TimeSpan.FromSeconds(10));
// Consumer with manual commit var consumerConfig = new ConsumerConfig { BootstrapServers = bootstrapServers, GroupId = groupId, EnableAutoCommit = false, IsolationLevel = IsolationLevel.ReadCommitted }; consumer = new ConsumerBuilder<string, string>(consumerConfig).Build();
processedOffsets = new HashSet<string>(); }
public void ProcessExactlyOnce(string topic, Func<object, object> handler) { consumer.Subscribe(topic);
try { while (true) { var result = consumer.Consume(TimeSpan.FromSeconds(1));
if (result == null) continue;
// Check if already processed string offsetKey = $"{result.Topic}-{result.Partition}-{result.Offset}";
if (processedOffsets.Contains(offsetKey)) { Console.WriteLine($"Skipping duplicate: {offsetKey}"); continue; }
try { // Begin transaction producer.BeginTransaction();
// Deserialize event var eventData = JsonSerializer.Deserialize<object>(result.Message.Value);
// Process message var processedResult = handler(eventData);
// Send result (in transaction) if (processedResult != null) { producer.Produce("output-topic", new Message<string, string> { Value = JsonSerializer.Serialize(processedResult) }); }
// Mark as processed processedOffsets.Add(offsetKey);
// Commit transaction producer.CommitTransaction(TimeSpan.FromSeconds(10)); } catch (Exception e) { // Abort transaction on error producer.AbortTransaction(TimeSpan.FromSeconds(10)); Console.Error.WriteLine($"Error processing: {e.Message}"); } } } catch (OperationCanceledException) { Console.WriteLine("Stopping processor..."); } finally { producer.Dispose(); consumer.Close(); } }}
// Usagevar processor = new ExactlyOnceProcessor("localhost:9092", "event-processors");
processor.ProcessExactlyOnce("input-topic", (eventData) => { Console.WriteLine($"Processing event: {eventData}"); // Process event... return new { processed = true, event = eventData };});Partitions Enable Parallelism
Partitions allow multiple consumers to process topic in parallel. Ordering guaranteed per partition.
Consumer Groups Scale
Consumer groups distribute partitions across consumers. Add consumers to scale throughput.
Offsets Track Progress
Offsets track consumer position. Commit offsets to resume after restart. Critical for reliability.
Exactly-Once is Complex
Exactly-once requires idempotent producer, transactions, and idempotent consumer. Use when duplicates are critical.