Repository navigation
feat: integrate Kafka event bus for task lifecycle events - #2
Merged
Merged
Conversation
EN: Apache Kafka 3.7.0 in KRaft combined mode (controller + broker), single-node cluster for development. Two client listeners: - INTERNAL://kafka:9092 for services inside the compose network - EXTERNAL://localhost:9094 for kcat/kafka-console-consumer from host Explicit single-node replication settings (offsets topic, transaction state log) since apache/kafka image assumes production defaults. Healthcheck via kafka-broker-api-versions.sh. task-service gains depends_on: kafka (service_healthy) as it's the only service producing/consuming events. Named volume kafka-data persists logs across compose down (not -v). ES: Broker Apache Kafka 3.7.0 en KRaft combined mode (controller + broker), cluster single-node para desarrollo. Dos listeners cliente: - INTERNAL://kafka:9092 para servicios dentro del compose network - EXTERNAL://localhost:9094 para kcat/kafka-console-consumer desde host Replication factor de topics internos declarado explícitamente por single-node (la imagen apache/kafka asume defaults de producción). Healthcheck vía kafka-broker-api-versions.sh. task-service añade depends_on: kafka (service_healthy) porque es el único servicio productor/consumer de eventos. Volumen kafka-data persiste los logs entre docker compose down (sin -v).
EN: Add spring-kafka dependency (version managed by Spring Boot BOM, resolved to 3.3.8 with kafka-clients 3.9.1). Configure Kafka in application.yml: - bootstrap-servers: localhost:9092 as base placeholder (explicit structure declared once, per ADR-005 principle) - producer.key-serializer: StringSerializer - producer.value-serializer: JsonSerializer (Spring Kafka) application-docker.yml overrides only the hostname to kafka:9092 (surgical override, no duplicated structure). Producer serializers apply across profiles; not duplicated in the docker profile since they don't change per environment. ES: Añadida dependencia spring-kafka (versión gestionada por el BOM de Spring Boot, resuelta a 3.3.8 con kafka-clients 3.9.1). Configuración Kafka en application.yml: - bootstrap-servers: localhost:9092 como placeholder base (estructura explícita declarada una sola vez, principio ADR-005) - producer.key-serializer: StringSerializer - producer.value-serializer: JsonSerializer (Spring Kafka) application-docker.yml sobrescribe solo el hostname a kafka:9092 (override quirúrgico, sin duplicación de estructura). Los serializers del productor aplican en todos los profiles; no se duplican en el docker profile porque no cambian entre entornos.
EN:
Introduce Kafka event publishing for task lifecycle in task-service.
New package com.mtole.task.kafka:
- KafkaConfig: NewTopic bean 'task.events' (3 partitions, replication 1).
Idempotent topic creation on startup, mirroring Flyway pattern:
infrastructure the app needs lives in the app's code, versioned,
auto-applied at startup.
- events.TaskEvent: thin event record (eventId, type, userId, taskId,
timestamp). Consumers requiring task details do lookup by taskId.
- events.TaskEventType: enum CREATED, UPDATED, STATUS_CHANGED, DELETED.
Mirrors the four ApplicationEvent types published by TaskService.
- TaskEventKafkaPublisher: four @TransactionalEventListener(AFTER_COMMIT)
methods listening to the existing local ApplicationEvents and
translating them to TaskEvent for Kafka publication with key=userId.
Design decisions (to be documented in ADR-007):
- Single topic per entity ('task.events') with type discriminator;
cross-topic ordering is impossible in Kafka.
- Key = userId aligns partitioning with ADR-004 identity axis and with
the /me/activity endpoint access pattern (compound index userId+ts).
- Thin events over fat events: payload as contract, state via lookup.
- eventId generated in publisher (Kafka contract, not local Spring event).
- AFTER_COMMIT phase guarantees no events published for failed
transactions. Does NOT cover broker failure mid-publish; transactional
outbox pattern documented as tech debt in ADR-007.
Verified empirically end-to-end: POST /tasks -> TaskCreatedEvent local
-> AFTER_COMMIT listener -> Kafka topic. All four event types observed
in kafka-console-consumer with key=userId and ordering preserved
within the partition.
ES:
Introduce la publicación de eventos de ciclo de vida de Task a Kafka
en task-service.
Nuevo package com.mtole.task.kafka:
- KafkaConfig: bean NewTopic 'task.events' (3 particiones, replication 1).
Creación idempotente del topic al arranque, espejo del patrón Flyway:
la infraestructura que la app necesita vive en el código de la app,
versionada, autoaplicada al arranque.
- events.TaskEvent: record de evento thin (eventId, type, userId,
taskId, timestamp). Consumers que necesiten detalles hacen lookup
por taskId.
- events.TaskEventType: enum CREATED, UPDATED, STATUS_CHANGED, DELETED.
Simetría con los cuatro ApplicationEvent locales publicados por
TaskService.
- TaskEventKafkaPublisher: cuatro métodos @TransactionalEventListener
(AFTER_COMMIT) que escuchan los ApplicationEvent locales existentes
y los traducen a TaskEvent para publicación Kafka con key=userId.
Decisiones de diseño (a documentar en ADR-007):
- Un topic por entidad ('task.events') con campo discriminador 'type';
el orden cross-topic es imposible en Kafka.
- Key = userId alinea el particionado con el eje de identidad de
ADR-004 y con el patrón de acceso de /me/activity (índice compuesto
userId+timestamp).
- Thin events por encima de fat events: payload como contrato, estado
vía lookup.
- eventId generado en el publisher (contrato Kafka, no del evento
local Spring).
- Fase AFTER_COMMIT garantiza que no se publican eventos de
transacciones fallidas. NO cubre fallo del broker mid-publish;
patrón transactional outbox documentado como deuda técnica en
ADR-007.
Verificado empíricamente end-to-end: POST /tasks -> TaskCreatedEvent
local -> listener AFTER_COMMIT -> topic Kafka. Los cuatro tipos de
evento observados en kafka-console-consumer con key=userId y orden
preservado dentro de la partición.
… listener retries Replace Spring Kafka's DefaultErrorHandler default (FixedBackOff(0ms, 9 attempts)) with ExponentialBackOff(initial=1s, multiplier=2, maxInterval=30s) and no total elapsed time cap. The default silently discards messages after exhausting retries, breaking the at-least-once guarantee documented in ADR-007 when downstream dependencies are unavailable beyond the retry budget. Empirically validated in session 15 (block 6.1): with the default handler, two of three events published during a Mongo outage were discarded (offset advanced without Mongo write). With this change, retries are effectively indefinite: consumer blocks on the failing partition until Mongo recovers. Preferred over silent data loss. Consumer lag becomes the visible signal of a stuck partition, replacing the invisible failure mode of the default configuration. --- fix(task-service): usar ExponentialBackOff sin techo total para reintentos del listener Kafka Sustituye el DefaultErrorHandler por defecto de Spring Kafka (FixedBackOff(0ms, 9 intentos)) por ExponentialBackOff(inicial=1s, multiplicador=2, maxInterval=30s) sin techo de tiempo total. El default descarta mensajes silenciosamente al agotar reintentos, rompiendo la garantía at-least-once documentada en ADR-007 cuando las dependencias downstream tardan más que el presupuesto de retry en recuperarse. Validado empíricamente en Sesión 15 (bloque 6.1): con el handler por defecto, dos de tres eventos publicados durante una caída de Mongo se descartaron (offset avanzó sin escritura en Mongo). Con este cambio los reintentos son indefinidos: el consumer se bloquea en la partición fallida hasta que Mongo se recupere. Preferible a la pérdida silenciosa. El lag del consumer pasa a ser la señal visible de una partición atascada, sustituyendo al fallo invisible de la configuración por defecto.
Adopt Apache Kafka as the event bus for domain lifecycle events in task-service. Documents the two-topic topology, key by userId, thin event payloads, at-least-once delivery with manual acknowledgement and indefinite exponential retry, and idempotency by _id = eventId. Consequences section captures the empirical finding from session 15: Spring Kafka's DefaultErrorHandler default silently discards messages after the retry budget is exhausted, breaking the at-least-once guarantee documented in D5. Corrected by commit 28da0d8. Includes a bidirectional idempotency observation: consumer group resets recover events historically lost while preventing duplicates for events already persisted. Alternatives considered: dead-letter topic, resource lookup at consume time, fat events, event enrichment for STATUS_CHANGED only, Repository.insert() vs MongoTemplate, and transactional outbox. Refs: 28da0d8 --- docs(adr): añadir ADR-007 sobre adopción del bus de eventos Kafka Adopta Apache Kafka como bus de eventos para eventos de ciclo de vida del dominio en task-service. Documenta la topología de dos topics, key por userId, payloads thin, entrega at-least-once con acknowledge manual y retry exponencial indefinido, e idempotencia por _id = eventId. La sección Consequences recoge el hallazgo empírico de Sesión 15: el DefaultErrorHandler por defecto de Spring Kafka descarta mensajes silenciosamente al agotar el presupuesto de retry, rompiendo la garantía at-least-once documentada en D5. Corregido en commit 28da0d8. Incluye una observación de idempotencia bidireccional: los reset del consumer group recuperan eventos históricamente perdidos mientras previenen duplicados de eventos ya persistidos. Alternativas consideradas: dead-letter topic, lookup de recurso al consumir, fat events, enriquecimiento solo para STATUS_CHANGED, Repository.insert() vs MongoTemplate, y transactional outbox. Refs: 28da0d8
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
EN
Introduces Kafka as event bus for task lifecycle events, in three
commits:
KRaft mode, dual client listeners (internal + external for kcat).
in application.yml, per-profile hostname override.
TaskEventType enum, TaskEventKafkaPublisher with
@TransactionalEventListener(AFTER_COMMIT).
Verified empirically end-to-end. All four event types (CREATED,
UPDATED, STATUS_CHANGED, DELETED) observed in kafka-console-consumer
with key=userId and ordering preserved within partition.
Pending in this branch (Session 14):
Draft PR to keep CI running while the feature grows.
ES
Introduce Kafka como event bus para eventos de ciclo de vida de Task,
en tres commits:
modo KRaft, dos listeners cliente (interno + externo para kcat).
serializers del productor en application.yml, override del
hostname por profile.
enum TaskEventType, TaskEventKafkaPublisher con
@TransactionalEventListener(AFTER_COMMIT).
Verificado empíricamente end-to-end. Los cuatro tipos (CREATED,
UPDATED, STATUS_CHANGED, DELETED) observados en kafka-console-consumer
con key=userId y orden preservado dentro de la partición.
Pendiente en esta rama (Sesión 14):
Draft PR para mantener CI corriendo mientras la feature crece.