From 56a39482e0d023dd276940157677ca24235833cf Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 10:57:20 +0900 Subject: [PATCH 1/3] test(cdc): fail when replica consumes its own DLT records --- .../replication/CdcReplicaConsumerTest.java | 21 +++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/replication/CdcReplicaConsumerTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/replication/CdcReplicaConsumerTest.java index c1adf451..815b7505 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/replication/CdcReplicaConsumerTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/replication/CdcReplicaConsumerTest.java @@ -6,9 +6,10 @@ import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.never; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -42,6 +43,22 @@ void delegatesSchemaChangesTopicToSchemaChangeApplier() { verify(processedDataReplicaApplier, never()).apply("xtrmetl-cdc.schema-changes", "k", "v"); } + @Test + void ignoresDeadLetterTopicSoRecoveryCannotCreateDltChains() { + CdcReplicaConsumer consumer = new CdcReplicaConsumer(processedDataReplicaApplier, schemaChangeReplicaApplier); + ConsumerRecord record = new ConsumerRecord<>( + "xtrmetl-cdc.public.processed_data.DLT", + 0, + 0L, + "k", + "poison" + ); + + consumer.onMessage(record); + + verifyNoInteractions(processedDataReplicaApplier, schemaChangeReplicaApplier); + } + @Test void delegatesNullTopicToProcessedDataApplier() { CdcReplicaConsumer consumer = new CdcReplicaConsumer(processedDataReplicaApplier, schemaChangeReplicaApplier); From af9b71d334f7cf265893b5286e83c67a79435048 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 11:00:38 +0900 Subject: [PATCH 2/3] refactor(cdc): name the terminal dead-letter topic suffix --- .../src/main/java/com/xtrmetl/cdc/replication/ReplicaTopics.java | 1 + 1 file changed, 1 insertion(+) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/replication/ReplicaTopics.java b/cdc-service/src/main/java/com/xtrmetl/cdc/replication/ReplicaTopics.java index 1a1e5b99..1828cbf9 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/replication/ReplicaTopics.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/replication/ReplicaTopics.java @@ -3,6 +3,7 @@ final class ReplicaTopics { static final String SCHEMA_CHANGES_SUFFIX = ".schema-changes"; + static final String DEAD_LETTER_SUFFIX = ".DLT"; private ReplicaTopics() { } From c88bf401f819a9cc7c9735746c4b55ae517b764f Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Mon, 10 Aug 2026 11:01:16 +0900 Subject: [PATCH 3/3] fix(cdc): keep dead-letter records out of replica appliers --- .../cdc/replication/CdcReplicaConsumer.java | 23 ++++++++++++++----- 1 file changed, 17 insertions(+), 6 deletions(-) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/replication/CdcReplicaConsumer.java b/cdc-service/src/main/java/com/xtrmetl/cdc/replication/CdcReplicaConsumer.java index 54e7d665..27086e4d 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/replication/CdcReplicaConsumer.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/replication/CdcReplicaConsumer.java @@ -5,6 +5,12 @@ import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; +/** + * Routes live CDC replica records to the matching replica applier. + * + *

Dead-letter topics are terminal quarantine inputs for operator recovery and are deliberately + * ignored by this live replica path, even when a configured topic pattern also matches them.

+ */ @Component @ConditionalOnProperty(prefix = "xtrmetl.replica", name = "enabled", havingValue = "true") public class CdcReplicaConsumer { @@ -13,10 +19,10 @@ public class CdcReplicaConsumer { private final SchemaChangeReplicaApplier schemaChangeReplicaApplier; /** - * CDC 복제를 처리하는 소비자 컴포넌트를 생성하고 필요한 복제 적용기(applier)를 주입한다. + * Creates the CDC replica consumer with the appliers used for live data and schema events. * - * @param processedDataReplicaApplier 처리된 데이터 복제 이벤트를 적용하는 객체 - * @param schemaChangeReplicaApplier 스키마 변경 복제 이벤트를 적용하는 객체 + * @param processedDataReplicaApplier applies supported data-row events to the replica + * @param schemaChangeReplicaApplier applies permitted schema-change events to the replica */ public CdcReplicaConsumer( ProcessedDataReplicaApplier processedDataReplicaApplier, @@ -27,11 +33,13 @@ public CdcReplicaConsumer( } /** - * 토픽 접미사에 따라 CDC 복제 메시지를 적절한 레플리카 applier로 라우팅하여 처리한다. + * Routes one live CDC record without feeding terminal dead-letter records back into appliers. * - *

레코드의 토픽이 ".schema-changes"로 끝나면 스키마 변경 처리기로 전달하고, 그렇지 않으면 처리된 데이터 복제기로 전달한다.

+ *

A topic ending in {@code .DLT} is a terminal dead-letter record and is acknowledged by the + * listener without invoking either replica applier. Schema-change topics are routed to the + * schema applier; every other live topic keeps the existing data-applier behavior.

* - * @param record Kafka로부터 수신된 레코드(토픽, 키, 값) + * @param record Kafka record containing the source topic, key, and Debezium-compatible value */ @KafkaListener( topicPattern = "${xtrmetl.replica.topic-pattern}", @@ -39,6 +47,9 @@ public CdcReplicaConsumer( ) public void onMessage(ConsumerRecord record) { String topic = record.topic(); + if (topic != null && topic.endsWith(ReplicaTopics.DEAD_LETTER_SUFFIX)) { + return; + } if (topic != null && topic.endsWith(ReplicaTopics.SCHEMA_CHANGES_SUFFIX)) { schemaChangeReplicaApplier.apply(topic, record.key(), record.value()); return;