fixed issue #726 , 优化SessionHandler的传输处理,提前序列化Entry
This commit is contained in:
@@ -92,11 +92,11 @@ public class EntryEventSink extends AbstractCanalEventSink<List<CanalEntry.Entry
|
||||
boolean hasHeartBeat = false;
|
||||
List<Event> events = new ArrayList<Event>();
|
||||
for (CanalEntry.Entry entry : entrys) {
|
||||
Event event = new Event(new LogIdentity(remoteAddress, -1L), entry);
|
||||
if (!doFilter(event)) {
|
||||
if (!doFilter(entry)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
Event event = new Event(new LogIdentity(remoteAddress, -1L), entry);
|
||||
events.add(event);
|
||||
hasRowData |= (entry.getEntryType() == EntryType.ROWDATA);
|
||||
hasHeartBeat |= (entry.getEntryType() == EntryType.HEARTBEAT);
|
||||
@@ -111,7 +111,7 @@ public class EntryEventSink extends AbstractCanalEventSink<List<CanalEntry.Entry
|
||||
} else {
|
||||
// 需要过滤的数据
|
||||
if (filterEmtryTransactionEntry && !CollectionUtils.isEmpty(events)) {
|
||||
long currentTimestamp = events.get(0).getEntry().getHeader().getExecuteTime();
|
||||
long currentTimestamp = events.get(0).getExecuteTime();
|
||||
// 基于一定的策略控制,放过空的事务头和尾,便于及时更新数据库位点,表明工作正常
|
||||
if (Math.abs(currentTimestamp - lastEmptyTransactionTimestamp) > emptyTransactionInterval
|
||||
|| lastEmptyTransactionCount.incrementAndGet() > emptyTransctionThresold) {
|
||||
@@ -126,15 +126,15 @@ public class EntryEventSink extends AbstractCanalEventSink<List<CanalEntry.Entry
|
||||
}
|
||||
}
|
||||
|
||||
protected boolean doFilter(Event event) {
|
||||
if (filter != null && event.getEntry().getEntryType() == EntryType.ROWDATA) {
|
||||
String name = getSchemaNameAndTableName(event.getEntry());
|
||||
protected boolean doFilter(CanalEntry.Entry entry) {
|
||||
if (filter != null && entry.getEntryType() == EntryType.ROWDATA) {
|
||||
String name = getSchemaNameAndTableName(entry);
|
||||
boolean need = filter.filter(name);
|
||||
if (!need) {
|
||||
logger.debug("filter name[{}] entry : {}:{}",
|
||||
name,
|
||||
event.getEntry().getHeader().getLogfileName(),
|
||||
event.getEntry().getHeader().getLogfileOffset());
|
||||
entry.getHeader().getLogfileName(),
|
||||
entry.getHeader().getLogfileOffset());
|
||||
}
|
||||
|
||||
return need;
|
||||
|
||||
+2
-2
@@ -18,7 +18,7 @@ public class HeartBeatEntryEventHandler extends AbstractCanalEventDownStreamHand
|
||||
public List<Event> before(List<Event> events) {
|
||||
boolean existHeartBeat = false;
|
||||
for (Event event : events) {
|
||||
if (event.getEntry().getEntryType() == EntryType.HEARTBEAT) {
|
||||
if (event.getEntryType() == EntryType.HEARTBEAT) {
|
||||
existHeartBeat = true;
|
||||
}
|
||||
}
|
||||
@@ -29,7 +29,7 @@ public class HeartBeatEntryEventHandler extends AbstractCanalEventDownStreamHand
|
||||
// 目前heartbeat和其他事件是分离的,保险一点还是做一下检查处理
|
||||
List<Event> result = new ArrayList<Event>();
|
||||
for (Event event : events) {
|
||||
if (event.getEntry().getEntryType() != EntryType.HEARTBEAT) {
|
||||
if (event.getEntryType() != EntryType.HEARTBEAT) {
|
||||
result.add(event);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -139,7 +139,7 @@ public class TimelineBarrier implements GroupBarrier<Event> {
|
||||
}
|
||||
|
||||
private Long getTimestamp(Event event) {
|
||||
return event.getEntry().getHeader().getExecuteTime();
|
||||
return event.getExecuteTime();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
+2
-2
@@ -113,11 +113,11 @@ public class TimelineTransactionBarrier extends TimelineBarrier {
|
||||
}
|
||||
|
||||
private boolean isTransactionBegin(Event event) {
|
||||
return event.getEntry().getEntryType() == EntryType.TRANSACTIONBEGIN;
|
||||
return event.getEntryType() == EntryType.TRANSACTIONBEGIN;
|
||||
}
|
||||
|
||||
private boolean isTransactionEnd(Event event) {
|
||||
return event.getEntry().getEntryType() == EntryType.TRANSACTIONEND;
|
||||
return event.getEntryType() == EntryType.TRANSACTIONEND;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -60,33 +60,33 @@ public class DummyEventStore implements CanalEventStore<Event> {
|
||||
}
|
||||
|
||||
public void put(Event data) throws InterruptedException, CanalStoreException {
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
}
|
||||
|
||||
public boolean put(Event data, long timeout, TimeUnit unit) throws InterruptedException, CanalStoreException {
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
return true;
|
||||
}
|
||||
|
||||
public boolean tryPut(Event data) throws CanalStoreException {
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
return true;
|
||||
}
|
||||
|
||||
public void put(List<Event> datas) throws InterruptedException, CanalStoreException {
|
||||
Event data = datas.get(0);
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
}
|
||||
|
||||
public boolean put(List<Event> datas, long timeout, TimeUnit unit) throws InterruptedException, CanalStoreException {
|
||||
Event data = datas.get(0);
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
return true;
|
||||
}
|
||||
|
||||
public boolean tryPut(List<Event> datas) throws CanalStoreException {
|
||||
Event data = datas.get(0);
|
||||
System.out.println("time:" + data.getEntry().getHeader().getExecuteTime());
|
||||
System.out.println("time:" + data.getExecuteTime());
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user