From 058e1fc546d4d2a7bd50c2a2e748c088b2cb5258 Mon Sep 17 00:00:00 2001 From: agapple Date: Tue, 28 May 2019 00:14:20 +0800 Subject: [PATCH] fixed delete message send mq bugfix --- .../otter/canal/common/MQMessageUtils.java | 31 +++++++------------ 1 file changed, 12 insertions(+), 19 deletions(-) diff --git a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java index d1598809..9161f8a5 100644 --- a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java +++ b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java @@ -222,33 +222,26 @@ public class MQMessageUtils { } else { for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { int hashCode = database.hashCode(); + CanalEntry.EventType eventType = rowChange.getEventType(); + List columns = null; + if (eventType == CanalEntry.EventType.DELETE) { + columns = rowData.getBeforeColumnsList(); + } else { + columns = rowData.getAfterColumnsList(); + } + if (hashMode.autoPkHash) { // isEmpty use default pkNames - for (CanalEntry.Column column : rowData.getAfterColumnsList()) { + for (CanalEntry.Column column : columns) { if (column.getIsKey()) { hashCode = hashCode ^ column.getValue().hashCode(); } } } else { - try { - CanalEntry.EventType eventType = RowChange.parseFrom(entry.getStoreValue()).getEventType(); - if(eventType == CanalEntry.EventType.DELETE){ - for (CanalEntry.Column column : rowData.getBeforeColumnsList()) { - if (checkPkNamesHasContain(hashMode.pkNames, column.getName())) { - hashCode = hashCode ^ column.getValue().hashCode(); - } - } + for (CanalEntry.Column column : columns) { + if (checkPkNamesHasContain(hashMode.pkNames, column.getName())) { + hashCode = hashCode ^ column.getValue().hashCode(); } - else { - for (CanalEntry.Column column : rowData.getAfterColumnsList()) { - if (checkPkNamesHasContain(hashMode.pkNames, column.getName())) { - hashCode = hashCode ^ column.getValue().hashCode(); - } - } - } - } - catch (InvalidProtocolBufferException e) { - e.printStackTrace(); } }