From 3e0d1df1f2f6bab48b62372720ec315f45f809c0 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Wed, 12 Aug 2026 21:25:16 +0900 Subject: [PATCH 1/3] test(cdc): replay malformed optional key RED on live develop --- .../spi/DebeziumChangeRecordMapperTest.java | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapperTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapperTest.java index 1f75c211..ffec5d11 100644 --- a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapperTest.java +++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapperTest.java @@ -80,6 +80,33 @@ void mapsDeleteUsingTopicWhenSourceMissing() { assertEquals(7, ((Number) record.getPk().get("id")).intValue()); } + @Test + void malformedKeyFallsBackToValidAfterIdentifierWithoutDroppingTheEvent() { + String value = """ + { + "payload": { + "op": "u", + "before": {"id": 41, "data": "old"}, + "after": {"id": 42, "data": "new"}, + "source": {"schema": "public", "table": "processed_data"} + } + } + """; + + Optional result = mapper.map( + "postgres-debezium", + "xtrmetl-cdc.public.processed_data", + "{malformed-key-json", + value + ); + + assertTrue(result.isPresent(), "a malformed optional key must not discard an otherwise valid CDC value"); + CanonicalChangeRecord record = result.get(); + assertEquals("u", record.getOp()); + assertEquals(42, ((Number) record.getPk().get("id")).intValue()); + assertEquals("new", record.getAfter().get("data")); + } + @Test void emptyValueReturnsEmpty() { assertTrue(mapper.map("s", "t", null, null).isEmpty()); From a6203fe688b9b1e7370f3d4eb7db0bc875957d43 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Wed, 12 Aug 2026 22:50:39 +0900 Subject: [PATCH 2/3] fix(cdc): isolate malformed optional Debezium keys --- .../xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java index 086e8cd5..ac4fe02b 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java @@ -105,13 +105,17 @@ public Optional map( } } - private Map extractPk(String keyJson) throws IOException { + private Map extractPk(String keyJson) { if (keyJson == null || keyJson.isBlank()) { return Map.of(); } - JsonNode root = objectMapper.readTree(keyJson); - JsonNode payload = root.has("payload") ? root.get("payload") : root; - return toMap(payload); + try { + JsonNode root = objectMapper.readTree(keyJson); + JsonNode payload = root.has("payload") ? root.get("payload") : root; + return toMap(payload); + } catch (IOException e) { + return Map.of(); + } } private static String[] schemaTableFromTopic(String topic) { From 2af6c9a7305ac8364f9e57691568136de6ab3c08 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Thu, 13 Aug 2026 05:10:17 +0900 Subject: [PATCH 3/3] docs(cdc): preserve malformed-key mapping contract --- .../xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java index ac4fe02b..df10e708 100644 --- a/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java +++ b/cdc-service/src/main/java/com/xtrmetl/cdc/spi/DebeziumChangeRecordMapper.java @@ -27,10 +27,17 @@ public DebeziumChangeRecordMapper(ObjectMapper objectMapper) { } /** + * Maps one Debezium value envelope and its optional key into a canonical change record. + * + *

A missing or malformed required value is rejected. A malformed optional key is treated + * as unavailable key metadata so a valid value can still supply the existing {@code id} + * fallback from its {@code after} or {@code before} object.

+ * * @param sourceId logical source id (e.g. {@code postgres-debezium}) - * @param topic Kafka / Debezium destination topic ({@code prefix.schema.table}) - * @param keyJson optional Debezium key JSON + * @param topic Kafka / Debezium destination topic ({@code prefix.schema.table}) + * @param keyJson optional Debezium key JSON * @param valueJson Debezium value JSON + * @return the mapped record when the required value envelope is valid, otherwise empty */ public Optional map( String sourceId,