Skip to main content
Un evento se mueve por Streaming Hub en tres etapas. El hub lo ingiere desde el stream interno, lo compara con tus suscripciones y lo despacha a cada destino que coincide. Esta página recorre cada etapa y las garantías de entrega que salen de ella.

Ingesta de eventos desde el stream


Streaming Hub consume el stream de la plataforma como mensajes CloudEvents 1.0 en modo de contenido binario. Los atributos de contexto de CloudEvents viajan en los headers del registro de Kafka (cada uno con el prefijo ce-). El cuerpo del evento es el valor del registro. El hub lee el enrutamiento y la identidad desde los headers sin deserializar el payload. El hub se suscribe solo a los temas de hechos de aplicación v3:
ce-source es un solo segmento en minúsculas y sin puntos. Usa letras, dígitos, guiones bajos o guiones. El tema no lleva sufijo de recurso, de evento ni de major de esquema. El hub despacha por los headers de CloudEvents. El conjunto de seguimiento excluye los temas de comandos (.commands) y los temas de dead-letter (.dlq). Los temas de hechos de aplicación de Midaz, Tracer, Matcher, Lender y Consignado ahora cumplen este contrato del conjunto de seguimiento. Un segundo consumidor lee los temas .dlq que escriben los productores. Ese consumidor sirve solo para observabilidad. La gramática del conjunto de seguimiento excluye un .dlq final, así que un registro de dead-letter nunca se convierte en una entrega a una de tus suscripciones. Los operadores leen esos registros con GET /admin/dlq. Consulta Análisis forense de la DLQ. Cada registro que el hub acepta debe llevar este conjunto de contexto: Para cada registro el hub devuelve uno de tres veredictos: El tenant admitido en la ingesta viene solo del atributo ce-tenantid validado que extrajo el parser. La ingesta no tiene un control de descarte por listado de tenants. Los eventos guardados coinciden solo con suscripciones que tengan el mismo id de tenant.
Tres cadenas relacionadas se parecen, pero son distintas. Nunca las confundas:
  • ce-type es el tipo de CloudEvents calificado por fuente: studio.lerian.<source>.<resource>.<event> (por ejemplo, studio.lerian.lender.loan_application.approved).
  • Tema de Kafka es un tema de hechos v3 por aplicación: lerian.streaming.<source> (por ejemplo, lerian.streaming.lender). El recurso, el evento y el major de esquema no aparecen en el tema.
  • X-Lerian-Event-Type es la clave simple <resource>.<event> en un webhook entregado (loan_application.approved). Omite el productor a propósito. Usa una fijación de origin en la suscripción cuando la misma clave puede venir de más de una aplicación.
Consulta el resumen de event streaming para el contrato de transmisión completo.

Deduplicación de consumo único


El stream interno es de al menos una vez, así que el mismo registro puede llegar más de una vez. Al entrar al inbox de eventos el hub deduplica por ce-id. Un id duplicado no escribe nada y cuenta como un descarte por deduplicación. El hub persiste los registros por partición de Kafka en una sola transacción. Confirma la posición del stream de una partición solo después de que la escritura de esa partición se confirma. Por eso una falla en una partición nunca pierde ni confirma dos veces los eventos de otra.

Coincidencia de suscripciones


Después de que el hub guarda un evento, resuelve cuáles de tus suscripciones deben recibirlo. La coincidencia de suscripciones es interna y no usa el catálogo: evalúa cada evento directamente sobre (tenant, origin, event type, schema major). La dimensión origin es opcional en la suscripción. Sin ella, la misma clave de evento de cualquier aplicación productora puede coincidir. Un evento que ninguna suscripción quiere produce cero entregas. Ese es un resultado correcto, no un error. La coincidencia admite una suscripción solo cuando está enabled y en verification_state = active. Estas son dos condiciones independientes, y la coincidencia requiere ambas. Consulta el modelo de suscripción.

Despacho a tus destinos


Por cada coincidencia, el hub crea un trabajo de entrega y lo pasa a un grupo de workers (ocho workers de forma predeterminada). Los workers reclaman trabajos vencidos de forma atómica con un lease corto y se alternan entre tenants, así que la acumulación de un solo tenant no deja sin recursos a los demás. Cada worker carga el payload guardado, descifra en memoria el secreto de firma o la credencial del destino, y hace exactamente un intento de entrega al sink. La entrega es al menos una vez: el hub puede entregar el mismo evento más de una vez (mediante reintentos o reentregas después de que un worker se reinicia). Cada entrega lleva un X-Lerian-Event-Id estable (el ce-id) para que deduplices por él, y un X-Lerian-Delivery-Id por intento que cambia en cada reintento. Trata un X-Lerian-Event-Id repetido como un duplicado y confírmalo sin volver a procesarlo. Consulta Consumo de eventos para el lado del consumidor de este contrato.

Reintentos y back-off


Cuando un intento falla, el hub programa el siguiente sobre una curva de back-off fija con jitter completo:
La curva abarca unas 28 horas en sus ocho pasos. Cada suscripción puede anularla con su propia programación (hasta 12 pasos, cada uno con un tope de 10 horas). Una anulación malformada vuelve a la curva predeterminada en lugar de hacer fallar la entrega.

Envío a dead-letter


El hub manda a dead-letter un evento que agota sus reintentos. Deja de intentarlo y registra el resultado final. El hub recupera un trabajo que, en cambio, sigue haciendo caer a un worker a mitad de intento y nunca registra un resultado. El límite de recuperación es cinco de forma predeterminada. Luego el hub lo manda a dead-letter como envenenado, así que un solo trabajo tóxico nunca puede ocupar un worker para siempre.

Circuit breaker


Cada par (tenant, destination) tiene un circuit breaker. Cuando un destino falla repetidamente, su circuit breaker se abre. Los siguientes trabajos para ese destino se descartan (se reprograman sin intento). Un solo endpoint roto no quema capacidad de workers ni golpea un destino que ya está sufriendo. El circuit breaker se recupera por sí solo una vez que el destino vuelve a tener éxito.

Desactivación automática de un destino roto


El hub desactiva automáticamente un destino que sigue roto durante mucho tiempo. La desactivación automática se dispara solo cuando el tramo de fallas abierto cumple las dos condiciones:
  • sostenido: dura al menos la ventana de fallas (120 horas de forma predeterminada), y
  • extendido: al menos 12 horas separan la primera falla de la última del tramo.
Ambos controles deben cumplirse, y una sola entrega exitosa borra el tramo, así que un problema breve nunca acumula hacia el veredicto. Cuando la desactivación automática se dispara, cambia la bandera enabled de la suscripción a false. Nunca toca verification_state. La entrega se detiene (la coincidencia requiere enabled), y el hub registra por qué. Recuperas una suscripción desactivada automáticamente con POST /v1/subscriptions/:id/verify, que vuelve a sondear el destino y, si tiene éxito, la vuelve a habilitar en su lugar. Un interruptor de apagado (STREAMING_HUB_AUTODISABLE_ENABLED) controla la desactivación automática, así que un operador puede desplegarla apagada. Consulta recuperar una suscripción desactivada automáticamente.

Garantías de ordenamiento


El stream interno está particionado por tenant. Los eventos de un tenant normalmente viajan en una sola partición, y el hub conserva su orden primero en entrar, primero en salir de extremo a extremo. El hub atribuye y deduplica por los headers de CloudEvents, no por la clave del registro de Kafka. Por eso, un productor que reparte (con salt) un tenant caliente entre varias particiones cambia solo la ubicación física. El hub sigue atribuyendo y deduplicando correctamente. El ordenamiento es la única concesión. Los eventos de un tenant con salt abarcan varias particiones sin garantía de orden entre particiones, así que ese tenant renuncia al FIFO estricto. La secuencia de llegada que asigna el hub refleja el orden en que los eventos llegaron al hub, no el orden en que se produjeron. Un tenant sin salt conserva el FIFO de una sola partición en todo el camino.

Próximos pasos


Gestión de suscripciones

El modelo de suscripción, los flujos de incorporación y la rotación de secretos.

Consumo de eventos

Verifica firmas, deduplica y lee eventos con pull.