| name | kafka-patterns |
| description | Apache Kafka patterns — producers, consumers, topics, consumer groups, exactly-once semantics, event sourcing. Use when working with kafka patterns. |
| domain | development |
| author | oyi77 |
| license | Apache-2.0 |
| subdomain | software-development |
| tags | ["coding","kafka","patterns","software-engineering","testing"] |
| version | 1.0.0 |
Overview
Apache Kafka is a distributed event streaming platform. This skill covers production-grade patterns for designing topics, building producers/consumers, managing consumer groups, implementing exactly-once semantics, and building event-sourced systems.
Capabilities
- Design topic hierarchies with partition strategies
- Build idempotent producers with exactly-once delivery
- Implement consumer groups with rebalance handling
- Use Schema Registry for Avro/Protobuf message evolution
- Build event sourcing and CQRS patterns
- Monitor lag, throughput, and consumer health
When to Use
Trigger phrases:
-
"kafka patterns"
-
"Apache Kafka patterns — producers, consumers, topics, consumer groups, exactly-o"
-
Building event-driven microservices
-
Need reliable message delivery with ordering guarantees
-
Implementing event sourcing or CQRS architectures
-
Processing high-throughput streaming data (millions of events/sec)
-
Decoupling producers and consumers in distributed systems
When NOT to Use
- Task is about deployment, not development (use deploy skills)
- Task is about code review, not writing (use review skills)
- You need to understand existing code first (use research skills)
- Task is about testing only (use test skills)
- Requirements are unclear (clarify first)
- Task is trivially simple (single line fix)
Pseudo Code
The kafka-patterns workflow follows a standard pipeline pattern.
Core flow:
# kafka-patterns primary flow
input = prepare(raw_data)
result = process(input, config={apache, consumer, consumers, event, exactly})
validate(result)
deliver(result)
Error handling:
on error:
log(error_details)
retry_with_backoff(max=3)
if still_failing: alert_and_escalate()
Topic Design
from confluent_kafka.admin import AdminClient, NewTopic
admin = AdminClient({'bootstrap.servers': 'localhost:9092'})
topic = NewTopic(
topic='orders.created',
num_partitions=12,
replication_factor=3,
config={
'retention.ms': str(7 * 24 * 60 * 60 * 1000),
'cleanup.policy': 'delete',
'min.insync.replicas': '2'
}
)
admin.create_topics([topic])
Idempotent Producer
from confluent_kafka import Producer
producer = Producer({
'bootstrap.servers': 'localhost:9092',
'enable.idempotence': True,
'acks': 'all',
'max.in.flight.requests': 5,
'retries': 2147483647,
'linger.ms': 5,
})
def delivery_callback(err, msg):
if err:
print(f'Delivery failed: {err}')
else:
print(f'Delivered to {msg.topic()}[{msg.partition()}]@{msg.offset()}')
producer.produce(
topic='orders.created',
key='user-12345',
value='{"order_id": "abc", "total": 99.99}',
callback=delivery_callback
)
producer.flush()
Consumer Group
from confluent_kafka import Consumer, KafkaError
consumer = Consumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'order-processor',
'auto.offset.reset': 'earliest',
'enable.auto.commit': False,
'isolation.level': 'read_committed',
})
consumer.subscribe(['orders.created'])
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())
try:
process_order(msg.value())
consumer.commit(msg)
except Exception as e:
handle_failure(msg, e)
Event Sourcing Pattern
events = {
'orders.events': {
'partitions': 12,
'key': 'order_id',
'schema': {
'event_type': 'string',
'aggregate_id': 'string',
'timestamp': 'long',
'payload': 'object',
'version': 'int'
}
}
}
def rebuild_order_state(order_id):
state = {}
for msg in consume_partition('orders.events', key=order_id):
event = deserialize(msg)
if event['event_type'] == 'OrderCreated':
state = event['payload']
elif event['event_type'] == 'OrderPaid':
state['paid'] = True
elif event['event_type'] == 'OrderShipped':
state['shipped'] = True
state['tracking'] = event['payload']['tracking']
return state
Schema Registry
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
schema_registry = SchemaRegistryClient({'url': 'http://localhost:8081'})
order_schema = '''
{
"type": "record",
"name": "Order",
"fields": [
{"name": "order_id", "type": "string"},
{"name": "user_id", "type": "string"},
{"name": "total", "type": "double"},
{"name": "items", "type": {"type": "array", "items": "string"}}
]
}
'''
serializer = AvroSerializer(schema_registry, order_schema)
Common Patterns
| Pattern | Use Case |
|---|
| Transactional outbox | Write DB + publish atomically |
| Event sourcing | Rebuild state from event log |
| CQRS | Separate read/write models |
| Dead letter queue | Handle poison messages |
| Compacted topics | Latest value per key (changelog) |
| Fan-out | Multiple consumers on same topic |
| Exactly-once | Idempotent producer + read_committed consumer |
How to Use
- Understand the requirement and existing codebase patterns
- Design the solution with error handling and testability in mind
- Implement incrementally with tests for each change
- Verify against expected outcomes (manual and automated)
- Document usage, edge cases, and integration points
- Review with team before merging to shared branches
Red Flags
- Skipping tests to ship faster: Untested code breaks in production when you least expect it
- No error handling in production code: Unhandled errors crash services and lose user data
- Hardcoded configuration values: Hardcoded values prevent environment switching and leak secrets
- Ignoring security implications: Missing input validation, auth bypasses, and injection vulnerabilities
- Over-engineering simple solutions: Premature abstraction adds complexity without proportional benefit
Verification
Process
- Analyze the task requirements
- Apply domain expertise
- Verify output quality
Anti-Rationalization Table
| Rationalization | Reality |
|---|
| "Tests slow me down" | Bugs slow you down 10x more. Tests are speed, not overhead. |
| "I will refactor later" | Technical debt compounds. Refactor as you go. |
| "It works on my machine" | If it is not in CI, it does not work. Ship proof, not claims. |