19 min de lectura

Objetivos de aprendizaje

Al final de este módulo, podrás:

  • Conectar productores y consumidores de Kafka a un clúster de IONOS CLOUD Event Streams for Apache Kafka utilizando autenticación mTLS y direcciones de brokers obtenidas de la salida de Terraform
  • Diseñar temas, particiones y grupos de consumidores para el flujo de eventos de TaskBoard, aplicando las restricciones de partición y replicación de IONOS CLOUD
  • Implementar entrega al menos una vez con confirmaciones manuales de offsets y un tema de letras muertas para el aislamiento de mensajes tóxicos
  • Configurar la confiabilidad del productor con idempotencia, `acks=all` y lotes, y ajustar la paralelización del consumidor frente al límite de número de particiones
  • Crear contratos de eventos con JSON Schema y evolucionarlos sin romper los consumidores existentes

Unidad 4.3: Integración de streaming de eventos

Introducción

Está conectando el servicio de API de TaskBoard con su worker de fondo. Cuando un usuario crea, mueve o cierra una tarea, la API no debe bloquearse al enviar un correo electrónico o un webhook. En su lugar, la API produce un evento task-changed hacia Kafka y devuelve la respuesta de inmediato, mientras que un servicio worker independiente consume esos eventos y dispara las notificaciones. Esto desacopla la ruta de la solicitud de los efectos secundarios lentos y le permite escalar el worker de forma independiente.

IONOS CLOUD Event Streams for Apache Kafka ejecuta Apache Kafka 4.0.0, por lo que sus bibliotecas de cliente de Kafka existentes funcionan sin cambios. Lo que difiere de un broker administrado por el usuario es la conexión: el plano de datos requiere TLS mutuo, no SASL ni texto plano, y las direcciones del broker y los certificados provienen del clúster que aprovisionó en la Unidad 2.5. Esta unidad comienza con la conexión del cliente y luego desarrolla el diseño de temas, la entrega confiable, el manejo de cartas muertas y la escalabilidad del consumidor, todo en el flujo de eventos de TaskBoard.

1. Conexión de clientes con mTLS

El plano de datos de Kafka de IONOS CLOUD está protegido con TLS mediante autenticación mutua. Ambas partes se autentican entre sí: su cliente verifica el certificado del broker contra la autoridad de certificación del clúster, y el broker verifica su certificado de cliente, que firma la propia autoridad de certificación del clúster. No existe un listener en texto plano ni un listener SASL/PLAIN. Cada productor y consumidor presenta un certificado de cliente.

Obtiene las credenciales, la autoridad de certificación, la clave privada del cliente y el certificado del cliente desde el punto de acceso del clúster. Estos son los mismos valores que el clúster expone a través de sus direcciones de broker y su API de acceso. Cada certificado de cliente es válido durante 365 días, por lo que la rotación de certificados debe incluirse en su manual de operaciones antes de que finalice el año.

1.1 Configuración de conexión del cliente

Una configuración de cliente de Kafka autogestionado utiliza security.protocol=SSL con un truststore (para la CA del clúster) y un keystore (para el certificado y la clave del cliente). La documentación de IONOS CLOUD muestra el archivo de propiedades del cliente canónico:

security.protocol=SSL
ssl.truststore.type=PKCS12
ssl.truststore.location=ca-cert.p12
ssl.truststore.password=changeit
ssl.endpoint.identification.algorithm=

ssl.keystore.type=PKCS12
ssl.keystore.location=admin.p12
ssl.keystore.password=adminp12pass

Nota: ssl.endpoint.identification.algorithm se deja vacío. Dado que el clúster utiliza una CA privada en lugar de certificados firmados públicamente, la verificación del nombre de host contra las direcciones del broker está deshabilitada aquí. El truststore valida la cadena en su lugar.

1.2 Conexión del productor en Python

La biblioteca confluent-kafka acepta las mismas claves SSL. Dirija el cliente a las direcciones del broker de su salida de Terraform y proporcione la CA, el certificado y la clave como archivos PEM.

from confluent_kafka import Producer

# Broker addresses come from the ionoscloud_kafka_cluster Terraform output.
conf = {
    "bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
    "security.protocol": "ssl",
    "ssl.ca.location": "ca-cert.pem",
    "ssl.certificate.location": "client-cert.pem",
    "ssl.key.location": "client-key.pem",
    "ssl.endpoint.identification.algorithm": "none",
    "client.id": "taskboard-api",
}

producer = Producer(conf)

El clúster ejecuta tres brokers en la topología del clúster, por lo que debe enumerar los tres en bootstrap.servers para garantizar la resiliencia de la conexión. El cliente solo necesita un broker accesible para iniciar la operación, pero enumerar los tres evita un punto único de falla durante el inicio.

2. Diseño de temas y particiones

El diseño de temas es donde las restricciones específicas de IONOS CLOUD se hacen sentir por primera vez. Utilice un tema por tipo de evento en lugar de un único tema de canalización masiva, ya que los consumidores se suscriben a los eventos que les interesan y puede configurar la retención y la partición según cada necesidad. Para TaskBoard, crea un tema task-changed para el trabajador y, más adelante, un tema de letra muerta task-changed.dlq.

2.1 Creación de temas mediante la API

La API de gestión de Kafka es regional y está separada de cloudapi. El host sigue el patrón https://kafka.<region>.ionos.com y, a diferencia del plano de datos, la API de gestión se autentica con un token Bearer, no con mTLS. Cree un tema con POST /clusters/{clusterId}/topics:

curl -X POST \
  'https://kafka.de-txl.ionos.com/clusters/e69b22a5-8fee-56b1-b6fb-4a07e4205ead/topics' \
  --header 'Content-Type: application/json' \
  --header 'Authorization: Bearer '"$IONOS_TOKEN" \
  --data '{
    "properties": {
      "name": "task-changed",
      "replicationFactor": 3,
      "numberOfPartitions": 6,
      "logRetention": { "retentionTime": 604800000 }
    }
  }'

Los parámetros del cuerpo son name (el único campo obligatorio), replicationFactor, numberOfPartitions y la configuración de logRetention. Lea los temas con GET /clusters/{clusterId}/topics y GET /clusters/{clusterId}/topics/{topicId}, que devuelven name, replicationFactor, numberOfPartitions y logRetention. La eliminación es DELETE /clusters/{clusterId}/topics/{topicId} y devuelve 202 Accepted, ya que la eliminación es asíncrona.

2.2 Restricciones de número de particiones y replicación

Elija el número de particiones como un múltiplo de 3 (3, 6, 9, 12, etc.) para evitar una distribución desigual de particiones entre los tres brokers. Un número de particiones que no sea un múltiplo del número de brokers deja que un broker asuma más particiones que los demás, lo que desequilibra la carga.

El factor de replicación recomendado es 3, coincidente con la topología de tres brokers, pero esta es una recomendación de IONOS CLOUD, no un valor predeterminado impuesto por el sistema; debe establecer replicationFactor explícitamente al crear un tema. Con un factor de replicación de 3, el clúster conserva tres copias de cada mensaje, por lo que el tema sobrevive a la pérdida de un broker sin pérdida de datos. La retención predeterminada del tema es de 604800000 ms (7 días). Establecer retentionTime en -1 no aplica ningún límite de tiempo, en cuyo caso la retención está limitada únicamente por el almacenamiento del clúster que haya aprovisionado (por ejemplo, 750 GB en el tamaño S, 1200 GB en M).

La siguiente tabla asocia los parámetros del cuerpo del tema con su función al diseñar un tema.

Parámetro Rol Valor de TaskBoard
name Identificador del tema (obligatorio) task-changed
replicationFactor Copias por partición entre brokers 3
numberOfPartitions Límite de paralelismo para consumidores 6 (múltiplo de 3)
logRetention.retentionTime Tiempo de vida del mensaje en ms (-1 = ilimitado) 604800000 (7 días)

Las particiones importan más allá del almacenamiento. El número de particiones es el límite estricto del paralelismo de consumidores dentro de un grupo de consumidores, por lo que dimensione este valor para el rendimiento máximo actual, ya que aumentar las particiones más adelante cambia el mapeo de clave a partición y altera las garantías de orden.

3. Fiabilidad del productor

Un productor con el enfoque de "disparar y olvidar" descarta mensajes cuando el broker experimenta interrupciones. Para los eventos de tareas de TaskBoard, donde la pérdida de un evento implica una notificación omitida, configure el productor para garantizar la durabilidad antes de optimizarlo para el rendimiento.

3.1 Confirmaciones, idempotencia y claves

Establezca acks=all para que el productor espere a que todas las réplicas en sincronía confirmen una escritura antes de considerarla exitosa. Combinado con un factor de replicación de 3, esto significa que un mensaje es durable entre brokers antes de que su código continúe. Habilite el productor idempotente para que las reintentos no creen duplicados en el lado del broker.

from confluent_kafka import Producer
import json

conf = {
    "bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
    "security.protocol": "ssl",
    "ssl.ca.location": "ca-cert.pem",
    "ssl.certificate.location": "client-cert.pem",
    "ssl.key.location": "client-key.pem",
    "ssl.endpoint.identification.algorithm": "none",
    "acks": "all",
    "enable.idempotence": True,
    "retries": 5,
    "linger.ms": 20,
    "batch.size": 65536,
}
producer = Producer(conf)

def on_delivery(err, msg):
    if err:
        # Persist for replay; never silently drop.
        print(f"delivery failed: {err}")
    else:
        print(f"delivered to {msg.topic()} [{msg.partition()}] @ {msg.offset()}")

def emit_task_changed(task_id: str, event: dict):
    producer.produce(
        topic="task-changed",
        key=task_id.encode("utf-8"),   # same task -> same partition -> ordered
        value=json.dumps(event).encode("utf-8"),
        on_delivery=on_delivery,
    )
    producer.poll(0)

La clave por task_id enruta todos los eventos de una tarea a la misma partición, lo que preserva el orden por tarea. Sin una clave, Kafka distribuye los registros de forma round-robin y un evento de "tarea cerrada" podría procesarse antes que "tarea creada".

3.2 Loteo y vaciado

linger.ms y batch.size intercambian unos pocos milisegundos de latencia por un rendimiento mucho mayor al agrupar registros por partición. Siempre ejecute flush() antes de que el proceso finalice, o los registros en búfer se perderán.

# At shutdown, block until all buffered messages are delivered or fail.
remaining = producer.flush(timeout=10)
if remaining > 0:
    raise RuntimeError(f"{remaining} messages not delivered before shutdown")

4. Grupos de consumidores, desplazamientos y semánticas de entrega

El trabajador de TaskBoard lee task-changed como un grupo de consumidores. Todos los consumidores del mismo grupo comparten las particiones, por lo que agregar pods de trabajadores aumenta el rendimiento hasta el número de particiones. Un cuarto consumidor en un tema de 3 particiones permanece inactivo.

4.1 Al menos una vez con confirmaciones manuales

Para la entrega de notificaciones, la semántica de al menos una vez es el valor predeterminado adecuado: procese el mensaje y luego confirme el desplazamiento. Si el trabajador se interrumpe después de procesar el mensaje pero antes de confirmar el desplazamiento, el mensaje se vuelve a entregar y la notificación podría activarse dos veces, lo cual se hace seguro mediante una clave de idempotencia. Lo contrario, confirmar antes de procesar (a lo sumo una vez), conlleva el riesgo de perder eventos en caso de una interrupción.

from confluent_kafka import Consumer, KafkaException
import json

conf = {
    "bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
    "security.protocol": "ssl",
    "ssl.ca.location": "ca-cert.pem",
    "ssl.certificate.location": "client-cert.pem",
    "ssl.key.location": "client-key.pem",
    "ssl.endpoint.identification.algorithm": "none",
    "group.id": "taskboard-notifier",
    "enable.auto.commit": False,        # manual commit = control over semantics
    "auto.offset.reset": "earliest",
}
consumer = Consumer(conf)
consumer.subscribe(["task-changed"])

while True:
    msg = consumer.poll(1.0)
    if msg is None:
        continue
    if msg.error():
        raise KafkaException(msg.error())
    event = json.loads(msg.value())
    send_notification(event)            # do the work first
    consumer.commit(msg)                # then commit only this offset

enable.auto.commit=False es la elección clave. El auto-compromiso avanza los desplazamientos mediante un temporizador, independientemente de si el procesamiento tuvo éxito o no, lo que provoca la pérdida silenciosa de mensajes cuyo procesamiento falló.

4.2 Escalado de consumidores y reequilibrio

La cantidad de particiones es el límite de paralelismo. Para escalar el trabajador de TaskBoard en Managed Kubernetes (véase la Unidad 3.2), aumente la cantidad de réplicas del Deployment hasta la cantidad de particiones, y cada nuevo pod se unirá al grupo taskboard-notifier y se le asignará un subconjunto de particiones a través de un reequilibrio. Durante un reequilibrio, las particiones dejan de consumirse brevemente, por lo que mantenga el procesamiento idempotente y los compromisos frecuentes para minimizar el reprocesamiento.

Particiones del tema Máximo de consumidores útiles Efecto de un consumidor adicional
3 3 El 4.º consumidor queda inactivo
6 6 Escala linealmente hasta 6
12 12 Margen para crecimiento futuro

Dimensione las particiones para su cantidad máxima de consumidores desde el inicio. Puede agregar particiones más adelante, pero al hacerlo se reorganiza la asignación de claves a particiones y se rompe el orden para las claves en proceso.

5. Manejo de mensajes en cola de errores y disciplina de esquema

No todos los mensajes pueden procesarse. Una carga de datos mal formada o una interrupción en un servicio aguas abajo seguirá fallando en cada reentrega, y un mensaje defectuoso que se reintente de forma indefinida bloqueará su partición completa.

5.1 Patrón de tema de cola de errores

Dirija los mensajes que agoten sus reintentos a un tema de cola de errores dedicado, task-changed.dlq, y luego confirme el desplazamiento original para que la partición principal continúe fluyendo. Inspeccione la cola de errores más adelante, corrija la causa y, de forma opcional, vuelva a reproducir los mensajes.

from confluent_kafka import Producer, Consumer
import json

dlq = Producer(conf_producer)  # same mTLS settings as the main producer

def process_with_dlq(consumer, msg, max_attempts=3):
    attempts = int(dict(msg.headers() or {}).get("x-attempts", b"0") or 0)
    try:
        send_notification(json.loads(msg.value()))
        consumer.commit(msg)
    except Exception as exc:
        if attempts + 1 >= max_attempts:
            dlq.produce(
                "task-changed.dlq",
                key=msg.key(),
                value=msg.value(),
                headers={"x-error": str(exc), "x-attempts": str(attempts + 1)},
            )
            dlq.flush(5)
            consumer.commit(msg)        # advance past the poison message
        else:
            raise                        # let it redeliver for a transient error

Distinga los fallos transitorios (un tiempo de espera en el servicio de notificaciones, que merece un reintento) de los permanentes (un evento mal formado, que va directamente a la DLQ). Reintentar ciegamente los fallos permanentes detiene la partición; enviar ciegamente a la DLQ los fallos transitorios provoca la pérdida de eventos recuperables.

5.2 Contratos de eventos con JSON Schema

Los productores y consumidores se implementan de forma independiente, por lo que la carga del evento es un contrato entre ellos. Defínala con JSON Schema y valide en ambos lados. Evolucione el esquema de forma aditiva: agregue campos opcionales, nunca elimine ni cambie el nombre de los existentes, para que los consumidores antiguos sigan funcionando mientras se implementan nuevos productores.

import jsonschema

TASK_CHANGED_V1 = {
    "type": "object",
    "required": ["task_id", "action", "occurred_at"],
    "properties": {
        "task_id": {"type": "string"},
        "action": {"enum": ["created", "moved", "closed"]},
        "occurred_at": {"type": "string", "format": "date-time"},
        "assignee": {"type": "string"},   # added in v1.1, optional = safe
    },
    "additionalProperties": False,
}

def validate_event(event: dict):
    jsonschema.validate(event, TASK_CHANGED_V1)  # raises ValidationError on breach

Incluya una versión del esquema en el evento o en un encabezado para que los consumidores puedan ramificar la lógica en función de dicha versión durante una ventana de migración. Esta disciplina es lo que le permite cambiar la forma de los eventos de TaskBoard sin realizar una reimplementación coordinada y sincronizada del productor y del trabajador.

Tarjeta rápida de referencia de la API

Puntos de acceso principales de la API de gestión de Kafka. El host es regional, por ejemplo https://kafka.de-txl.ionos.com:

Método Punto de acceso Descripción
POST /clusters Crear un clúster de Kafka
GET /clusters/{clusterId} Obtener detalles del clúster y direcciones de los brokers
POST /clusters/{clusterId}/topics Crear un tema
GET /clusters/{clusterId}/topics Listar todos los temas
DELETE /clusters/{clusterId}/topics/{topicId} Eliminar un tema (devuelve 202)
GET /clusters/{clusterId}/users/{userId}/access Obtener las credenciales de acceso mTLS (CA, certificado, clave)

Autenticación de la API de gestión: Authorization: Bearer <token> Autenticación del plano de datos (productores/consumidores): TLS mutuo con certificado del cliente

Laboratorio de código

Objetivo: Producir y consumir eventos de TaskBoard task-changed en un clúster Kafka de IONOS CLOUD mediante un grupo de consumidores, confirmaciones manuales de desplazamientos y enrutamiento a una cola de mensajes no procesables en caso de fallo.

Requisitos previos:

  • Cuenta de IONOS CLOUD con token de API (IONOS_TOKEN)
  • Un clúster Kafka aprovisionado desde la Unidad 2.5 con direcciones de brokers disponibles
  • Python 3.10 o superior con confluent-kafka y jsonschema instalados
  • El certificado CA del clúster, el certificado del cliente y la clave del cliente guardados como archivos PEM

Paso 1: Obtener las direcciones de los brokers y las credenciales de acceso

CLUSTER_ID="e69b22a5-8fee-56b1-b6fb-4a07e4205ead"
curl -s -X GET \
  "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID" \
  -H "Authorization: Bearer $IONOS_TOKEN" | python3 -m json.tool

Salida esperada:

{
  "id": "e69b22a5-8fee-56b1-b6fb-4a07e4205ead",
  "properties": { "name": "...", "version": "4.0.0", "size": "S",
    "connections": [ { "brokerAddresses": ["192.168.1.101/24", ...] } ] }
}

Paso 2: Crear los temas principales y de mensajes no procesados

for T in task-changed task-changed.dlq; do
  curl -s -X POST "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics" \
    -H "Content-Type: application/json" -H "Authorization: Bearer $IONOS_TOKEN" \
    --data "{\"properties\":{\"name\":\"$T\",\"replicationFactor\":3,\"numberOfPartitions\":6}}"
done

Salida esperada:

{"id":"ae085c4c-...","type":"topic","properties":{"name":"task-changed",
 "replicationFactor":3,"numberOfPartitions":6}}

Paso 3: Generar un evento de cambio de tarea

producer.produce("task-changed", key=b"task-42",
    value=b'{"task_id":"task-42","action":"created","occurred_at":"2026-06-05T10:00:00Z"}')
producer.flush(10)
print("produced")

Salida esperada:

delivered to task-changed [3] @ 0
produced

Paso 4: Consumir con un grupo de consumidores y confirmación manual

consumer.subscribe(["task-changed"])
msg = consumer.poll(5.0)
print(msg.topic(), msg.partition(), msg.offset(), msg.value())
consumer.commit(msg)

Salida esperada:

task-changed 3 0 b'{"task_id":"task-42",...}'

Paso 5: Disparar una carta muerta con una carga de datos no válida

producer.produce("task-changed", key=b"task-99", value=b'{not valid json')
producer.flush(10)
# Consumer's process_with_dlq routes it to task-changed.dlq after max_attempts

Salida esperada:

delivered to task-changed.dlq [1] @ 0

Paso 6: Verificar que el tema de mensajes no procesados recibió el mensaje

dlq_consumer.subscribe(["task-changed.dlq"])
m = dlq_consumer.poll(5.0)
print(dict(m.headers()), m.value())

Salida esperada:

{'x-error': b'...', 'x-attempts': b'3'} b'{not valid json'

Lista de verificación:

  • [ ] Temas task-changed y task-changed.dlq creados con un factor de replicación de 3
  • [ ] El productor envía con acks=all y la idempotencia habilitada
  • [ ] El consumidor confirma los desplazamientos manualmente después del procesamiento
  • [ ] Un mensaje mal formado se dirige al tema de mensajes muertos, no a un bucle de reintentos infinito

Limpieza:

for ID in $(curl -s "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics" \
  -H "Authorization: Bearer $IONOS_TOKEN" | python3 -c \
  "import sys,json;[print(t['id']) for t in json.load(sys.stdin)['items']]"); do
  curl -s -X DELETE "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics/$ID" \
    -H "Authorization: Bearer $IONOS_TOKEN"
done

Errores comunes

  1. Uso de SASL o texto plano en lugar de mTLS

    • Problema: El cliente se queda colgado al conectar o falla con SSL handshake failed, aunque las credenciales parezcan correctas.
    • Por qué ocurre: El plano de datos de IONOS CLOUD solo acepta TLS mutuo. No existe un listener de SASL/PLAIN ni de texto plano, y los brokers utilizan una CA privada, por lo que un truststore predeterminado los rechaza.
    • Solución: Establezca security.protocol=ssl, proporcione la CA del clúster como ssl.ca.location, el certificado y la clave del cliente, y desactive la verificación del nombre de host con ssl.endpoint.identification.algorithm=none.
  2. Auto-commit que descarta silenciosamente mensajes fallidos

    • Problema: Algunas notificaciones de tareas nunca se ejecutan, pero no aparecen errores en los registros y el retardo del consumidor es cero.
    • Por qué ocurre: Con enable.auto.commit=True, los offsets avanzan mediante un temporizador sin importar si send_notification tuvo éxito, por lo que un mensaje que generó una excepción se marca como consumido.
    • Solución: Establezca enable.auto.commit=False y llame a consumer.commit(msg) solo después de que el procesamiento tenga éxito, dirigiendo los fallos permanentes al tema de mensajes no procesados.
  3. Número de particiones que no es múltiplo del número de brokers

    • Problema: Un broker muestra una carga de CPU y disco más alta que los otros dos, y el rendimiento se estanca por debajo de lo esperado.
    • Por qué ocurre: Con 3 brokers, un tema creado con 4 o 5 particiones se distribuye de forma desigual, por lo que un broker aloja un líder de partición adicional.
    • Solución: Cree temas con un número de particiones que sea múltiplo de 3 (3, 6, 9, 12) para que las particiones se distribuyan de manera uniforme entre los tres brokers.

Resumen

Ahora puede integrar IONOS CLOUD Event Streams for Apache Kafka en el código de la aplicación como un componente de primera clase de la arquitectura de TaskBoard. Conecta productores y consumidores mediante TLS mutuo utilizando las direcciones de los brokers y los certificados de su aprovisionado clúster, crea temas a través de la API de gestión regional con autenticación Bearer y diseña configuraciones de particiones y réplicas que respeten la topología de tres brokers. Además, cuenta con un productor confiable con idempotencia y acks=all, un grupo de consumidores con confirmaciones manuales de entrega al menos una vez, un tema de mensajes no procesables (dead-letter) para mensajes defectuosos y contratos de JSON Schema que permiten que el productor y el consumidor evolucionen de forma independiente.

Puntos clave:

  • El plano de datos de Kafka requiere TLS mutuo con un certificado de cliente; la API de gestión utiliza tokens Bearer. Son dos modelos de autenticación diferentes en dos puntos finales diferentes.
  • La cantidad de particiones de un tema debe ser múltiplo de 3 para distribuirlas de manera uniforme entre los tres brokers; el factor de réplica recomendado es 3, pero debe configurarse de forma explícita, ya que no es un valor predeterminado impuesto por el sistema.
  • La entrega al menos una vez se logra deshabilitando la confirmación automática y confirmando los desplazamientos (offsets) solo después de un procesamiento exitoso; debe hacer que el procesamiento sea idempotente para tolerar la reentrega.
  • La cantidad de particiones es el límite máximo fijo para la paralelización de los consumidores; dimensione las particiones para la carga máxima antes de desplegar, ya que agregarlas posteriormente interrumpe el orden de las claves.
  • Los temas de mensajes no procesables (dead-letter) aíslan los mensajes defectuosos para que una carga defectuosa no bloquee su partición, y JSON Schema aplica el contrato entre el productor y el consumidor.

Terminología importante:

  • mTLS (TLS mutuo): Tanto el cliente como el broker se autentican mutuamente mediante certificados; es el único método de autenticación del plano de datos de IONOS CLOUD, con certificados de cliente válidos por 365 días.
  • Grupo de consumidores: Un conjunto de consumidores que comparten el trabajo de las particiones de un tema; es la unidad de escalado horizontal para el trabajador de TaskBoard.
  • Confirmación de desplazamiento (offset commit): Registro de la posición hasta la que un consumidor ha procesado; confirmar después del procesamiento produce una entrega al menos una vez.
  • Tema de mensajes no procesables (DLQ): Un tema separado para mensajes que fallan en el procesamiento después de reintentos, manteniendo la partición principal desbloqueada.
  • Factor de réplica: El número de copias de cada partición entre brokers; en IONOS CLOUD es 3, coincidiendo con el clúster de tres brokers, de modo que el tema sobrevive a la pérdida de un broker.

Próximos pasos

Siga aprendiendo: Unidad 4.4: Integración de AI Model Hub

Temas relacionados: