Repository navigation
fix(flow): grupo de consumidor próprio na fila de gatilhos de jornada do ClickHouse (CRM-777) - #128
Conversation
…ornada A tabela Kafka journey_trigger_kafka_queue nomeava o grupo temporal-workers, o mesmo dos workers que consomem os gatilhos. O ClickHouse só produz nesse tópico, mas um consumidor iniciado nessa tabela dividiria as partições com os workers e eles perderiam gatilhos. - A fila de gatilhos passa a usar <KAFKA_GROUP_ID>-journey-triggers-clickhouse. - A guarda de boot (ensureKafkaEngineBroker) compara também o grupo e recria a tabela quando ele diverge, para a troca chegar a tabelas que já existem. A fila de contact_events passa o grupo dela também. Co-Authored-By: Claude Code <[EMAIL_REDACTED]>
Reviewer's GuideA implementação separa defensivamente o consumer group da tabela ClickHouse de gatilhos de jornada dos workers temporais e estende a guarda de boot para detectar divergências de grupo no DDL existente, recriando a MV e a tabela quando necessário; os testes cobrem parsing, recriação e wiring das filas. Sequence diagram for ClickHouse Kafka queue boot guardsequenceDiagram
participant ClickHouseService
participant ClickHouse
participant Kafka
ClickHouseService->>ClickHouse: ensureKafkaEngineBroker(expectedBrokers, expectedGroup)
ClickHouse-->>ClickHouseService: Existing Kafka Engine DDL
ClickHouseService->>ClickHouseService: extractKafkaBrokers()
ClickHouseService->>ClickHouseService: extractKafkaGroup()
alt broker or consumer group differs
ClickHouseService->>ClickHouse: DROP VIEW IF EXISTS events_to_journey_triggers_mv
ClickHouseService->>ClickHouse: DROP TABLE IF EXISTS journey_trigger_kafka_queue
ClickHouseService->>ClickHouse: CREATE TABLE journey_trigger_kafka_queue
ClickHouse->>Kafka: Produce to journey-triggers
else broker and group match
ClickHouseService-->>ClickHouse: Keep existing table
end
File-Level Changes
Tips and commandsInteracting with Sourcery
Customizing Your ExperienceAccess your dashboard to:
Getting Help
|
There was a problem hiding this comment.
Hey - I've reviewed your changes and they look great!
Sourcery assessment
Needs a human reviewer. The change alters the Kafka consumer group used by the ClickHouse engine table, which can change partition ownership and cause journey-trigger messages to be consumed by the wrong component or skipped. Reverting restores the configuration, but messages already consumed or offsets already advanced may not be recoverable by the revert alone.
Summary
A tabela Kafka
journey_trigger_kafka_queuedo ClickHouse nomeava o consumer grouptemporal-workers, o mesmo dos workers que consomem os gatilhos de jornada.<KAFKA_GROUP_ID>-journey-triggers-clickhouse.events_to_journey_triggers_mv. Então hoje ele não entra no grupo, e isso foi conferido emsystem.kafka_consumers.ensureKafkaEngineBroker) passa a comparar também o grupo. Quando ele diverge, a guarda derruba a MV e a tabela para recriá-las, como já fazia com o broker. Sem isso, a troca não chegaria às tabelas que já existem.expectedGroupé opcional.contact_eventspassa o grupo que já tinha, sem efeito.Security
contact_eventsnesse intervalo de milissegundos não viram gatilho.Test plan
npx jest src/modules/processing/clickhouse/→ 18/18. Os casos novos:temporal-workers.npx jest --ci --maxWorkers=2→ 1134 passed, 0 falhas.npm run typecheckenpm run buildok.temporal-workers;contact_eventssem passar o grupo;uses consumer group 'temporal-workers' … Recreating itaparece, e a tabela é recriada com o grupo novo.already points …, sem recriação.contact_eventschegou ao tópicojourney-triggers; o offset passou de 390 para 391.Changed Files
src/modules/processing/clickhouse/clickhouse.service.tssrc/modules/processing/clickhouse/clickhouse.service.spec.tssrc/modules/processing/clickhouse/clickhouse.service.contact-events-broker.spec.tsRelated PRs
Linked Issue
🤖 Generated with Claude Code