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..a7aa821c 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
@@ -22,15 +22,28 @@ public class DebeziumChangeRecordMapper {
private final ObjectMapper objectMapper;
+ /**
+ * Creates a mapper using the caller-provided JSON codec.
+ *
+ * @param objectMapper mapper used for Debezium key and value envelopes
+ */
public DebeziumChangeRecordMapper(ObjectMapper objectMapper) {
this.objectMapper = objectMapper;
}
/**
+ * Maps one Debezium value envelope and its optional key metadata.
+ *
+ *
The value envelope is authoritative and must be valid JSON. The key envelope is optional
+ * metadata: when it is malformed, mapping continues and derives a primary key from the value's
+ * {@code after.id} or {@code before.id} field when available. This prevents malformed optional
+ * metadata from discarding an otherwise valid change event.
+ *
* @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 mapped canonical record, or empty when the required value envelope is blank or malformed
*/
public Optional map(
String sourceId,
@@ -83,7 +96,7 @@ public Optional map(
Map afterMap = toMap(after);
Map beforeMap = toMap(before);
- Map pk = extractPk(keyJson);
+ Map pk = extractOptionalPk(keyJson);
if (pk.isEmpty() && afterMap.containsKey("id")) {
pk = Map.of("id", afterMap.get("id"));
} else if (pk.isEmpty() && beforeMap.containsKey("id")) {
@@ -105,13 +118,17 @@ public Optional map(
}
}
- private Map extractPk(String keyJson) throws IOException {
+ private Map extractOptionalPk(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 ignored) {
+ return Map.of();
+ }
}
private static String[] schemaTableFromTopic(String topic) {
diff --git a/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumMalformedKeyFallbackTest.java b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumMalformedKeyFallbackTest.java
new file mode 100644
index 00000000..e74ee15d
--- /dev/null
+++ b/cdc-service/src/test/java/com/xtrmetl/cdc/spi/DebeziumMalformedKeyFallbackTest.java
@@ -0,0 +1,35 @@
+package com.xtrmetl.cdc.spi;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/** Fail-first regression for optional malformed Debezium key metadata. */
+class DebeziumMalformedKeyFallbackTest {
+
+ @Test
+ void malformedOptionalKeyFallsBackToValidAfterIdentifier() {
+ DebeziumChangeRecordMapper mapper = new DebeziumChangeRecordMapper(new ObjectMapper());
+ 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(), "malformed optional key metadata must not discard a valid value envelope");
+ CanonicalChangeRecord record = result.orElseThrow();
+ assertEquals(42, ((Number) record.getPk().get("id")).intValue());
+ assertEquals("new", record.getAfter().get("data"));
+ }
+}