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..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,
@@ -105,13 +112,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) {
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());