diff --git a/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java b/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java index dfa4280b..d8821f46 100644 --- a/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java +++ b/protocol/src/main/java/com/alibaba/otter/canal/protocol/FlatMessage.java @@ -149,11 +149,19 @@ public class FlatMessage implements Serializable { } List flatMessages = new ArrayList<>(); + List entrys = null; + if (message.isRaw()) { + List rawEntries = message.getRawEntries(); + entrys = new ArrayList(rawEntries.size()); + for (ByteString byteString : rawEntries) { + CanalEntry.Entry entry = CanalEntry.Entry.parseFrom(byteString); + entrys.add(entry); + } + } else { + entrys = message.getEntries(); + } - List rawEntries = message.getRawEntries(); - - for (ByteString byteString : rawEntries) { - CanalEntry.Entry entry = CanalEntry.Entry.parseFrom(byteString); + for (CanalEntry.Entry entry : entrys) { if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN || entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) { continue; diff --git a/protocol/src/main/java/com/alibaba/otter/canal/protocol/Message.java b/protocol/src/main/java/com/alibaba/otter/canal/protocol/Message.java index 35adf15d..751e09d0 100644 --- a/protocol/src/main/java/com/alibaba/otter/canal/protocol/Message.java +++ b/protocol/src/main/java/com/alibaba/otter/canal/protocol/Message.java @@ -19,7 +19,6 @@ public class Message implements Serializable { private static final long serialVersionUID = 1234034768477580009L; private long id; - @Deprecated private List entries = new ArrayList(); // row data for performance, see: // https://github.com/alibaba/canal/issues/726 diff --git a/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java b/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java index a4f74c90..b07062ee 100644 --- a/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java +++ b/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java @@ -1,11 +1,12 @@ package com.alibaba.otter.canal.server.embedded; -import java.util.*; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.ServiceLoader; import java.util.concurrent.TimeUnit; -import com.alibaba.otter.canal.spi.CanalMetricsProvider; -import com.alibaba.otter.canal.spi.CanalMetricsService; -import com.alibaba.otter.canal.spi.NopCanalMetricsService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.slf4j.MDC; @@ -14,6 +15,7 @@ import org.springframework.util.CollectionUtils; import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; import com.alibaba.otter.canal.instance.core.CanalInstance; import com.alibaba.otter.canal.instance.core.CanalInstanceGenerator; +import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.ClientIdentity; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.position.LogPosition; @@ -22,7 +24,11 @@ import com.alibaba.otter.canal.protocol.position.PositionRange; import com.alibaba.otter.canal.server.CanalServer; import com.alibaba.otter.canal.server.CanalService; import com.alibaba.otter.canal.server.exception.CanalServerException; +import com.alibaba.otter.canal.spi.CanalMetricsProvider; +import com.alibaba.otter.canal.spi.CanalMetricsService; +import com.alibaba.otter.canal.spi.NopCanalMetricsService; import com.alibaba.otter.canal.store.CanalEventStore; +import com.alibaba.otter.canal.store.memory.MemoryEventStoreWithBuffer; import com.alibaba.otter.canal.store.model.Event; import com.alibaba.otter.canal.store.model.Events; import com.google.common.base.Function; @@ -40,12 +46,12 @@ import com.google.protobuf.ByteString; */ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements CanalServer, CanalService { - private static final Logger logger = LoggerFactory.getLogger(CanalServerWithEmbedded.class); + private static final Logger logger = LoggerFactory.getLogger(CanalServerWithEmbedded.class); private Map canalInstances; // private Map lastRollbackPostions; private CanalInstanceGenerator canalInstanceGenerator; private int metricsPort; - private CanalMetricsService metrics = NopCanalMetricsService.NOP; + private CanalMetricsService metrics = NopCanalMetricsService.NOP; private static class SingletonHolder { @@ -207,7 +213,7 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C * b. 如果timeout不为null * 1. timeout为0,则采用get阻塞方式,获取数据,不设置超时,直到有足够的batchSize数据才返回 * 2. timeout不为0,则采用get+timeout方式,获取数据,超时还没有batchSize足够的数据,有多少返回多少 - * + * * 注意: meta获取和数据的获取需要保证顺序性,优先拿到meta的,一定也会是优先拿到数据,所以需要加同步. (不能出现先拿到meta,拿到第二批数据,这样就会导致数据顺序性出现问题) * */ @@ -239,12 +245,23 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C } else { // 记录到流式信息 Long batchId = canalInstance.getMetaManager().addBatch(clientIdentity, events.getPositionRange()); - List entrys = Lists.transform(events.getEvents(), new Function() { + boolean raw = isRaw(canalInstance.getEventStore()); + List entrys = null; + if (raw) { + entrys = Lists.transform(events.getEvents(), new Function() { - public ByteString apply(Event input) { - return input.getRawEntry(); - } - }); + public ByteString apply(Event input) { + return input.getRawEntry(); + } + }); + } else { + entrys = Lists.transform(events.getEvents(), new Function() { + + public CanalEntry.Entry apply(Event input) { + return input.getEntry(); + } + }); + } if (logger.isInfoEnabled()) { logger.info("get successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]", clientIdentity.getClientId(), @@ -255,7 +272,7 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C } // 直接提交ack ack(clientIdentity, batchId); - return new Message(batchId, true, entrys); + return new Message(batchId, raw, entrys); } } } @@ -283,7 +300,7 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C * b. 如果timeout不为null * 1. timeout为0,则采用get阻塞方式,获取数据,不设置超时,直到有足够的batchSize数据才返回 * 2. timeout不为0,则采用get+timeout方式,获取数据,超时还没有batchSize足够的数据,有多少返回多少 - * + * * 注意: meta获取和数据的获取需要保证顺序性,优先拿到meta的,一定也会是优先拿到数据,所以需要加同步. (不能出现先拿到meta,拿到第二批数据,这样就会导致数据顺序性出现问题) * */ @@ -311,7 +328,8 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C } if (CollectionUtils.isEmpty(events.getEvents())) { - // logger.debug("getWithoutAck successfully, clientId:{} batchSize:{} but result + // logger.debug("getWithoutAck successfully, clientId:{} + // batchSize:{} but result // is null", // clientIdentity.getClientId(), // batchSize); @@ -319,12 +337,23 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C } else { // 记录到流式信息 Long batchId = canalInstance.getMetaManager().addBatch(clientIdentity, events.getPositionRange()); - List entrys = Lists.transform(events.getEvents(), new Function() { + boolean raw = isRaw(canalInstance.getEventStore()); + List entrys = null; + if (raw) { + entrys = Lists.transform(events.getEvents(), new Function() { - public ByteString apply(Event input) { - return input.getRawEntry(); - } - }); + public ByteString apply(Event input) { + return input.getRawEntry(); + } + }); + } else { + entrys = Lists.transform(events.getEvents(), new Function() { + + public CanalEntry.Entry apply(Event input) { + return input.getEntry(); + } + }); + } if (logger.isInfoEnabled()) { logger.info("getWithoutAck successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]", clientIdentity.getClientId(), @@ -333,7 +362,7 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C batchId, events.getPositionRange()); } - return new Message(batchId, true, entrys); + return new Message(batchId, raw, entrys); } } @@ -515,17 +544,25 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C // 发现provider, 进行初始化 if (list.size() > 1) { logger.warn("Found more than one CanalMetricsProvider, use the first one."); - //报告冲突 + // 报告冲突 for (CanalMetricsProvider p : list) { logger.warn("Found CanalMetricsProvider: {}.", p.getClass().getName()); } } - //默认使用第一个 + // 默认使用第一个 CanalMetricsProvider provider = list.get(0); this.metrics = provider.getService(); } } + private boolean isRaw(CanalEventStore eventStore) { + if (eventStore instanceof MemoryEventStoreWithBuffer) { + return ((MemoryEventStoreWithBuffer) eventStore).isRaw(); + } + + return true; + } + // ========= setter ========== public void setCanalInstanceGenerator(CanalInstanceGenerator canalInstanceGenerator) { diff --git a/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java b/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java index 76124ec4..d34450ed 100644 --- a/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java +++ b/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java @@ -20,6 +20,7 @@ import com.alibaba.otter.canal.sink.CanalEventDownStreamHandler; import com.alibaba.otter.canal.sink.CanalEventSink; import com.alibaba.otter.canal.sink.exception.CanalSinkException; import com.alibaba.otter.canal.store.CanalEventStore; +import com.alibaba.otter.canal.store.memory.MemoryEventStoreWithBuffer; import com.alibaba.otter.canal.store.model.Event; /** @@ -42,7 +43,8 @@ public class EntryEventSink extends AbstractCanalEventSink props = entry.getHeader().getPropsList(); if (props != null) { @@ -60,6 +64,16 @@ public class Event implements Serializable { } } } + + if (raw) { + // build raw + this.rawEntry = entry.toByteString(); + this.rawLength = rawEntry.size(); + } else { + this.entry = entry; + // 按照3倍的event length预估 + this.rawLength = entry.getHeader().getEventLength() * 3; + } } public LogIdentity getLogIdentity() { @@ -150,6 +164,14 @@ public class Event implements Serializable { this.rowsCount = rowsCount; } + public CanalEntry.Entry getEntry() { + return entry; + } + + public void setEntry(CanalEntry.Entry entry) { + this.entry = entry; + } + public String toString() { return ToStringBuilder.reflectionToString(this, CanalToStringStyle.DEFAULT_STYLE); }