πŸ“‘ CapΓ­tulo 10 Β· Nivel Avanzado

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.

⏱️ Lectura: ~65 min
🎯 Nivel: Avanzado
πŸ“Œ Kafka: 3.7+ (KRaft mode)

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.

graph TD subgraph "Producers" P1[πŸ›’ Orders Service] P2[πŸ‘€ User Service] P3[πŸ“Š Analytics App] end subgraph "Kafka Cluster (KRaft Mode - Kafka 3.3+)" B1[πŸ–₯️ Broker 1
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

πŸ†• KRaft Mode (Kafka sin Zookeeper)

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)