Apache Kafka
Apache Kafka es la plataforma de event streaming distribuida mΓ‘s usada del mundo. Es el corazΓ³n de las arquitecturas modernas event-driven y el elemento clave para datos en tiempo real.
Arquitectura de Kafka
Kafka es un distributed commit log: una secuencia ordenada, inmutable y replicada de mensajes. Los productores escriben al final del log, los consumidores leen desde cualquier posiciΓ³n.
Leader: orders-0, users-1] B2[π₯οΈ Broker 2
Leader: orders-1, users-0] B3[π₯οΈ Broker 3
Leader: orders-2, users-2] B1 <--> B2 B2 <--> B3 B1 <--> B3 end subgraph "Topics" T1[π orders
3 partitions, RF=3] T2[π user-events
6 partitions, RF=3] end subgraph "Consumer Groups" CG1[βοΈ analytics-group
Consumer 1, 2, 3] CG2[βοΈ notifications-group
Consumer A, B] CG3[βοΈ data-lake-group
Kafka Connect Sink] end P1 --> T1 P2 --> T2 P3 --> T1 T1 --> B1 T1 --> B2 T2 --> B2 T2 --> B3 B1 --> CG1 B2 --> CG2 B3 --> CG3 style B1 fill:#7c0d0d,stroke:#e34234,color:#fff style B2 fill:#7c0d0d,stroke:#e34234,color:#fff style B3 fill:#7c0d0d,stroke:#e34234,color:#fff
Conceptos Fundamentales
- Broker: Servidor Kafka. Almacena topics y sirve a producers/consumers. Un cluster tiene mΓΊltiples brokers para HA y escalabilidad.
- Topic: Canal lΓ³gico de mensajes. Como una tabla de base de datos pero para eventos. Immutable y ordenado por particiΓ³n.
- Partition: SubdivisiΓ³n fΓsica de un topic. El paralelismo de Kafka = nΓΊmero de particiones. Los mensajes dentro de una particiΓ³n estΓ‘n ordenados.
- Offset: PosiciΓ³n de un mensaje dentro de una particiΓ³n. Incremental, inmutable. Los consumers guardan su offset para saber dΓ³nde van.
- Replication Factor: NΓΊmero de copias de cada particiΓ³n. RF=3 significa que cada particiΓ³n estΓ‘ en 3 brokers. Tolerancia a fallos de RF-1 brokers.
- Leader/Follower: Cada particiΓ³n tiene un broker Leader (recibe escrituras) y N Followers (replican). Si el Leader falla, un Follower se convierte en Leader.
Kafka 3.3+ introdujo el modo KRaft (Kafka Raft) que elimina la dependencia de Apache Zookeeper. En 2026, KRaft es el modo por defecto y recomendado. Simplifica el despliegue (un sistema menos que operar) y mejora la escalabilidad (hasta millones de particiones vs ~200K con Zookeeper).
Topics y Particiones
# ββ GESTIΓN DE TOPICS βββββββββββββββββββββββββββββββββββββββββββββββββββββ
from confluent_kafka.admin import AdminClient, NewTopic, ConfigResource
admin = AdminClient({'bootstrap.servers': 'kafka-1:9092,kafka-2:9092'})
# Crear topic con configuraciones
topics_to_create = [
NewTopic(
topic='orders',
num_partitions=12, # Paralelismo = max 12 consumers simultΓ‘neos
replication_factor=3, # 3 copias para HA
config={
'retention.ms': str(7 * 24 * 60 * 60 * 1000), # 7 dΓas
'retention.bytes': str(50 * 1024 * 1024 * 1024), # 50GB por particiΓ³n
'min.insync.replicas': '2', # Al menos 2 rΓ©plicas deben confirmar el write
'compression.type': 'snappy',
'max.message.bytes': str(1 * 1024 * 1024), # 1MB mΓ‘ximo por mensaje
'cleanup.policy': 'delete', # o 'compact' para change-data-capture
}
),
NewTopic(
topic='user-profiles',
num_partitions=6,
replication_factor=3,
config={
'cleanup.policy': 'compact', # Log compaction: mantiene ΓΊltimo valor por key
'min.cleanable.dirty.ratio': '0.1', # Compactar cuando 10% estΓ‘ "sucio"
}
)
]
result = admin.create_topics(topics_to_create)
for topic, future in result.items():
try:
future.result()
print(f"β
Topic {topic} created")
except Exception as e:
print(f"β Failed to create {topic}: {e}")
# ββ DECISIΓN: ΒΏCuΓ‘ntas particiones? βββββββββββββββββββββββββββββββββββββββ
# Regla general:
# - Max throughput esperado / throughput por particiΓ³n (tΓpico ~10-100 MB/s)
# - Si esperas 100 consumers mΓ‘ximo β 100 particiones
# - Para empezar: 12-24 particiones por topic de producciΓ³n
# - Puedes AUMENTAR particiones despuΓ©s, pero no REDUCIR
# - Demasiadas particiones = mayor overhead de metadata en el broker
# ββ LOG COMPACTION: Mantener el ΓΊltimo estado de cada key βββββββββββββββββ
# Γtil para change data capture y tablas de lookup
# Topic con cleanup.policy=compact:
# DespuΓ©s de compactaciΓ³n, solo queda el ΓLTIMO mensaje por key
# [key=user1, action=create] β [key=user1, action=update] β compacted β [key=user1, action=update]
# [key=user2, value=null] β DELETE marker: elimina todos los mensajes de key=user2
Producers
"""
Kafka Producer: enviar mensajes a Kafka
confluent-kafka-python (librdkafka bajo el capΓ³ - la mΓ‘s eficiente)
"""
import json
from confluent_kafka import Producer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
# ββ PRODUCER BΓSICO ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
producer_conf = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092,kafka-3:9092',
# CONFIABILIDAD
'acks': 'all', # Espera confirmaciΓ³n de TODOS los ISR (mΓ‘s seguro)
# 'acks': '1' # Solo leader (mΓ‘s rΓ‘pido, puede perder datos si leader falla)
# 'acks': '0' # Fire-and-forget (mΓ‘s rΓ‘pido, puede perder datos)
'enable.idempotence': True, # Garantiza exactly-once delivery
'retries': 5, # Reintentos en caso de error transitorio
'retry.backoff.ms': 100,
# PERFORMANCE
'linger.ms': 10, # Esperar hasta 10ms para batch de mensajes
'batch.size': 65536, # 64KB de batch mΓ‘ximo (aumenta throughput)
'compression.type': 'snappy', # CompresiΓ³n en el producer
'buffer.memory': 33554432, # 32MB de buffer en memoria
}
producer = Producer(producer_conf)
def delivery_report(err, msg):
"""Callback llamado cuando el mensaje es entregado o falla."""
if err is not None:
print(f'β Message delivery failed: {err}')
else:
print(f'β
Message delivered to {msg.topic()}[{msg.partition()}] @ offset {msg.offset()}')
# ββ ENVIAR MENSAJES ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def send_order_event(order: dict) -> None:
try:
producer.produce(
topic='orders',
key=str(order['customer_id']), # Key determina la particiΓ³n
value=json.dumps(order).encode('utf-8'),
headers={
'event-type': b'order_placed',
'source-service': b'orders-api',
'schema-version': b'v2',
},
on_delivery=delivery_report
)
producer.poll(0) # Servir callbacks sin bloquear
except BufferError:
# El buffer interno estΓ‘ lleno, esperar y reintentar
producer.flush(timeout=10)
producer.produce(topic='orders', key=str(order['customer_id']),
value=json.dumps(order).encode('utf-8'))
# Enviar batch de Γ³rdenes
orders = [{"order_id": i, "customer_id": i % 100, "amount": i * 9.99} for i in range(1000)]
for order in orders:
send_order_event(order)
# IMPORTANTE: flush() antes de terminar para enviar mensajes pendientes en el buffer
producer.flush(timeout=30)
print(f"Enviados {len(orders)} eventos de orden")
# ββ PARTICIΓN KEY STRATEGY βββββββββββββββββββββββββββββββββββββββββββββββββ
# Si KEY es None β Round-robin entre particiones (mΓ‘ximo throughput, sin orden)
# Si KEY = customer_id β Todos los eventos del mismo cliente van a la misma particiΓ³n
# β Orden garantizado por cliente (ΓΊtil para state machines por usuario)
# Si KEY = order_id β DistribuciΓ³n uniforme, orden por pedido
Consumers y Consumer Groups
"""
Kafka Consumer y Consumer Groups
El paralelismo = nΓΊmero de particiones activas = nΓΊmero de consumers activos
"""
from confluent_kafka import Consumer, KafkaException, TopicPartition
import signal
import sys
consumer_conf = {
'bootstrap.servers': 'kafka-1:9092,kafka-2:9092',
'group.id': 'analytics-pipeline', # Identificador del Consumer Group
# AUTO OFFSET MANAGEMENT (simple pero menos control)
'auto.offset.reset': 'earliest', # latest (solo nuevos) o earliest (desde inicio)
'enable.auto.commit': True, # Hace commit automΓ‘ticamente
'auto.commit.interval.ms': 5000, # Cada 5 segundos
# PERFORMANCE
'fetch.min.bytes': 1, # Min bytes para iniciar un fetch
'fetch.max.wait.ms': 500, # Max tiempo de espera para fetch.min.bytes
'max.poll.records': 500, # Max mensajes por poll()
}
consumer = Consumer(consumer_conf)
consumer.subscribe(['orders', 'returns']) # Subscribirse a mΓΊltiples topics
# ββ CONSUMER LOOP ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
running = True
def signal_handler(sig, frame):
global running
running = False
signal.signal(signal.SIGINT, signal_handler)
messages_processed = 0
batch = []
try:
while running:
# poll(): recibe mensajes (bloqueante hasta timeout)
msg = consumer.poll(timeout=1.0)
if msg is None:
# Timeout: no hay mensajes nuevos
if batch:
process_batch(batch) # Procesar el batch acumulado
batch = []
continue
if msg.error():
raise KafkaException(msg.error())
# Deserializar mensaje
event = json.loads(msg.value().decode('utf-8'))
event['_kafka_topic'] = msg.topic()
event['_kafka_partition'] = msg.partition()
event['_kafka_offset'] = msg.offset()
batch.append(event)
messages_processed += 1
# Procesar en batches de 100
if len(batch) >= 100:
process_batch(batch)
batch = []
except Exception as e:
print(f"Error: {e}")
finally:
if batch:
process_batch(batch)
consumer.close() # Deja que otro consumer tome sus particiones
print(f"Procesados: {messages_processed:,} mensajes")
# ββ COMMIT MANUAL: mΓ‘s control ββββββββββββββββββββββββββββββββββββββββββββ
consumer_manual_conf = {**consumer_conf, 'enable.auto.commit': False}
consumer_manual = Consumer(consumer_manual_conf)
consumer_manual.subscribe(['orders'])
while True:
msg = consumer_manual.poll(timeout=1.0)
if msg is None:
continue
try:
event = json.loads(msg.value())
process_event(event) # Procesar el mensaje
consumer_manual.commit(msg) # Solo hacer commit si el proceso fue exitoso
except Exception as e:
# No hacer commit β el mensaje serΓ‘ re-procesado
print(f"Error procesando {msg.offset()}: {e}")
# ββ CONSUMER GROUP: Particiones se asignan automΓ‘ticamente βββββββββββββββββ
# Topic: orders (12 particiones)
# Consumer Group: analytics-pipeline (3 consumers)
# Resultado: 4 particiones por consumer
#
# Si se agrega consumer 4: rebalance β 3 particiones para algunos, 2 para otros
# Si consumer 3 cae: rebalance β sus 4 particiones se redistribuyen entre los restantes
#
# Lag = (Latest Offset - Committed Offset)
# Lag = 0 β Consumer estΓ‘ al dΓa
# Lag >> 0 β Consumer estΓ‘ atrasado (escalar consumers o optimizar procesamiento)
Kafka Streams
"""
Kafka Streams con Python (Faust - el mejor wrapper Python para Kafka Streams)
Para casos simples donde no quieres Flink completo
"""
import faust
app = faust.App('order-processor',
broker='kafka://kafka-1:9092',
value_serializer='json')
# Definir los modelos de datos
class Order(faust.Record, serializer='json'):
order_id: int
customer_id: int
amount: float
status: str
class OrderStats(faust.Record, serializer='json'):
customer_id: int
total_orders: int
total_revenue: float
avg_order: float
# Topics
orders_topic = app.topic('orders', value_type=Order)
order_stats_topic = app.topic('order-stats', value_type=OrderStats)
# ββ STREAM PROCESSING: Transformar y enriquecer eventos ββββββββββββββββββββ
@app.agent(orders_topic)
async def process_orders(stream):
async for order in stream:
# Filtrar Γ³rdenes no vΓ‘lidas
if order.amount <= 0:
continue
# Enriquecer con datos externos (por ej, info del cliente desde cache)
enriched = {
'order_id': order.order_id,
'customer_id': order.customer_id,
'amount': order.amount,
'status': order.status,
'revenue_category': 'large' if order.amount >= 500 else 'small',
}
# Re-publicar al topic enriquecido
await order_stats_topic.send(
key=str(order.customer_id),
value=enriched
)
# ββ TABLA (KTABLE): AgregaciΓ³n stateful ββββββββββββββββββββββββββββββββββββ
customer_stats_table = app.Table('customer-stats', default=lambda: {'orders': 0, 'revenue': 0.0})
@app.agent(orders_topic)
async def aggregate_customer_stats(stream):
async for order in stream:
# Actualizar el estado en la tabla (backed por Kafka topic changelog)
stats = customer_stats_table[order.customer_id]
stats['orders'] += 1
stats['revenue'] += order.amount
customer_stats_table[order.customer_id] = stats
# Publicar el estado actualizado
await order_stats_topic.send(
key=str(order.customer_id),
value=OrderStats(
customer_id=order.customer_id,
total_orders=stats['orders'],
total_revenue=stats['revenue'],
avg_order=stats['revenue'] / stats['orders']
)
)
if __name__ == '__main__':
app.main()
Schema Registry
El Schema Registry gestiona los schemas de los mensajes Kafka (Avro, JSON Schema, Protobuf), garantizando compatibilidad entre producers y consumers cuando el schema evoluciona.
"""
Confluent Schema Registry con Avro
Garantiza que producers y consumers hablan el mismo "idioma"
"""
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer, AvroDeserializer
from confluent_kafka.serialization import StringSerializer, SerializationContext, MessageField
from confluent_kafka import Producer, Consumer
# Schema del evento
order_avro_schema = """
{
"type": "record",
"name": "Order",
"namespace": "com.company.orders",
"doc": "Represents a customer order event",
"fields": [
{"name": "order_id", "type": "long"},
{"name": "customer_id", "type": "int"},
{"name": "amount", "type": "double"},
{"name": "status", "type": "string"},
{"name": "created_at", "type": {"type": "long", "logicalType": "timestamp-millis"}},
{"name": "currency", "type": "string", "default": "USD"},
{"name": "notes", "type": ["null", "string"], "default": null}
]
}
"""
schema_registry_conf = {'url': 'http://schema-registry:8081'}
schema_registry_client = SchemaRegistryClient(schema_registry_conf)
# ββ PRODUCER CON AVRO ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
avro_serializer = AvroSerializer(schema_registry_client, order_avro_schema)
key_serializer = StringSerializer('utf_8')
producer = Producer({'bootstrap.servers': 'kafka:9092'})
order_data = {
'order_id': 1001,
'customer_id': 501,
'amount': 299.99,
'status': 'pending',
'created_at': int(datetime.now().timestamp() * 1000),
'currency': 'USD',
'notes': None
}
producer.produce(
topic='orders',
key=key_serializer(str(order_data['customer_id']), SerializationContext('orders', MessageField.KEY)),
value=avro_serializer(order_data, SerializationContext('orders', MessageField.VALUE)),
on_delivery=delivery_report
)
# ββ COMPATIBILIDAD DE SCHEMAS ββββββββββββββββββββββββββββββββββββββββββββββββ
# BACKWARD: nuevos consumers leen mensajes de producers antiguos (default)
# FORWARD: consumers antiguos leen mensajes de nuevos producers
# FULL: ambas direcciones (mΓ‘s restrictivo)
# Ver y cambiar la compatibilidad
schema_registry_client.set_compatibility(subject_name='orders-value', level='BACKWARD')
# Ver los schemas registrados
schemas = schema_registry_client.get_versions('orders-value')
print(f"Versions: {schemas}")
# EvoluciΓ³n de schema compatible (BACKWARD):
# β
Agregar campo con default
# β
Cambiar default de un campo
# β Eliminar campo sin default
# β Cambiar tipo de campo (int β string)
Kafka Connect
Kafka Connect es el framework para integrar Kafka con sistemas externos (bases de datos, storage, APIs) sin escribir cΓ³digo de producer/consumer.
# ββ KAFKA CONNECT: CDC desde PostgreSQL β Kafka ββββββββββββββββββββββββββββ
# Debezium: el conector CDC mΓ‘s popular
# Captura cambios de PostgreSQL (INSERT, UPDATE, DELETE) como eventos Kafka
# 1. Activar replication en PostgreSQL
# postgresql.conf: wal_level = logical
# pg_hba.conf: host replication debezium 0.0.0.0/0 md5
# 2. Configurar el conector Debezium via REST API
import requests
debezium_config = {
"name": "postgres-orders-cdc",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "debezium",
"database.password": "debezium",
"database.dbname": "production",
"database.server.name": "production",
"table.include.list": "public.orders,public.order_items",
# Snapshot inicial de la tabla completa
"snapshot.mode": "initial",
# Formato de salida
"key.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"key.converter.schema.registry.url": "http://schema-registry:8081",
"value.converter.schema.registry.url": "http://schema-registry:8081",
# Topics de output: production.public.orders, production.public.order_items
"topic.prefix": "production",
}
}
response = requests.post(
"http://kafka-connect:8083/connectors",
json=debezium_config
)
print(f"Connector created: {response.status_code}")
# 3. Verificar el conector
response = requests.get("http://kafka-connect:8083/connectors/postgres-orders-cdc/status")
print(response.json())
# Cada cambio en la tabla 'orders' genera un evento en el topic 'production.public.orders':
# {
# "before": null, # null para INSERT
# "after": {"order_id": 1001, "amount": 299.99, "status": "pending"},
# "source": {"db": "production", "table": "orders", "ts_ms": 1718000000000},
# "op": "c", # c=create, u=update, d=delete, r=read (snapshot)
# }
# ββ SINK CONNECTOR: Kafka β S3 (para el Data Lake) βββββββββββββββββββββββββ
s3_sink_config = {
"name": "s3-data-lake-sink",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "4",
"topics": "production.public.orders",
"s3.region": "us-east-1",
"s3.bucket.name": "my-data-lake",
"s3.part.size": "67108864", # 64MB parts
"flush.size": "1000", # Flush cada 1000 mensajes
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
"parquet.codec": "snappy",
"locale": "en_US",
"timezone": "UTC",
"timestamp.extractor": "RecordField",
"timestamp.field": "created_at",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"path.format": "'year'=YYYY/'month'=MM/'day'=dd/'hour'=HH",
"partition.duration.ms": "3600000", # 1 hora por particiΓ³n
}
}
requests.post("http://kafka-connect:8083/connectors", json=s3_sink_config)