public class EntryEventSink extends AbstractCanalEventSink<List<CanalEntry.Entry>> implements CanalEventSink<List<CanalEntry.Entry>> {
private static final Logger logger = LoggerFactory.getLogger(EntryEventSink.class); private static final int maxFullTimes = 10; private CanalEventStore<Event> eventStore; protected boolean filterTransactionEntry = false; protected boolean filterEmtryTransactionEntry = true; protected long emptyTransactionInterval = 5 * 1000; protected long emptyTransctionThresold = 8192;
protected volatile long lastTransactionTimestamp = 0L; protected AtomicLong lastTransactionCount = new AtomicLong(0L); protected volatile long lastEmptyTransactionTimestamp = 0L; protected AtomicLong lastEmptyTransactionCount = new AtomicLong(0L); protected AtomicLong eventsSinkBlockingTime = new AtomicLong(0L); protected boolean raw;
public EntryEventSink(){ addHandler(new HeartBeatEntryEventHandler()); }
public void start() { super.start(); Assert.notNull(eventStore);
if (eventStore instanceof MemoryEventStoreWithBuffer) { this.raw = ((MemoryEventStoreWithBuffer) eventStore).isRaw(); }
for (CanalEventDownStreamHandler handler : getHandlers()) { if (!handler.isStart()) { handler.start(); } } }
public void stop() { super.stop();
for (CanalEventDownStreamHandler handler : getHandlers()) { if (handler.isStart()) { handler.stop(); } } }
public boolean filter(List<Entry> event, InetSocketAddress remoteAddress, String destination) {
return false; }
public boolean sink(List<CanalEntry.Entry> entrys, InetSocketAddress remoteAddress, String destination) throws CanalSinkException, InterruptedException { return sinkData(entrys, remoteAddress); }
private boolean sinkData(List<CanalEntry.Entry> entrys, InetSocketAddress remoteAddress) throws InterruptedException { boolean hasRowData = false; boolean hasHeartBeat = false; List<Event> events = new ArrayList<>(); for (CanalEntry.Entry entry : entrys) { if (!doFilter(entry)) { continue; }
if (filterTransactionEntry && (entry.getEntryType() == EntryType.TRANSACTIONBEGIN || entry.getEntryType() == EntryType.TRANSACTIONEND)) { long currentTimestamp = entry.getHeader().getExecuteTime(); if (lastTransactionCount.incrementAndGet() <= emptyTransctionThresold && Math.abs(currentTimestamp - lastTransactionTimestamp) <= emptyTransactionInterval) { continue; } else { if (entry.getEntryType() == EntryType.TRANSACTIONEND) { lastTransactionCount.set(0L); lastTransactionTimestamp = currentTimestamp; } } }
hasRowData |= (entry.getEntryType() == EntryType.ROWDATA); hasHeartBeat |= (entry.getEntryType() == EntryType.HEARTBEAT); Event event = new Event(new LogIdentity(remoteAddress, -1L), entry, raw); events.add(event); }
if (hasRowData || hasHeartBeat) { return doSink(events); } else { if (filterEmtryTransactionEntry && !CollectionUtils.isEmpty(events)) { long currentTimestamp = events.get(0).getExecuteTime(); if (Math.abs(currentTimestamp - lastEmptyTransactionTimestamp) > emptyTransactionInterval || lastEmptyTransactionCount.incrementAndGet() > emptyTransctionThresold) { lastEmptyTransactionCount.set(0L); lastEmptyTransactionTimestamp = currentTimestamp; return doSink(events); } }
return true; } }
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, entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset()); }
return need; } else { return true; } }
protected boolean doSink(List<Event> events) { for (CanalEventDownStreamHandler<List<Event>> handler : getHandlers()) { events = handler.before(events); } long blockingStart = 0L; int fullTimes = 0; do { if (eventStore.tryPut(events)) { if (fullTimes > 0) { eventsSinkBlockingTime.addAndGet(System.nanoTime() - blockingStart); } for (CanalEventDownStreamHandler<List<Event>> handler : getHandlers()) { events = handler.after(events); } return true; } else { if (fullTimes == 0) { blockingStart = System.nanoTime(); } applyWait(++fullTimes); if (fullTimes % 100 == 0) { long nextStart = System.nanoTime(); eventsSinkBlockingTime.addAndGet(nextStart - blockingStart); blockingStart = nextStart; } }
for (CanalEventDownStreamHandler<List<Event>> handler : getHandlers()) { events = handler.retry(events); }
} while (running && !Thread.interrupted()); return false; }
private void applyWait(int fullTimes) { int newFullTimes = fullTimes > maxFullTimes ? maxFullTimes : fullTimes; if (fullTimes <= 3) { Thread.yield(); } else { LockSupport.parkNanos(1000 * 1000L * newFullTimes); }
}
private String getSchemaNameAndTableName(CanalEntry.Entry entry) { return entry.getHeader().getSchemaName() + "." + entry.getHeader().getTableName(); }
public void setEventStore(CanalEventStore<Event> eventStore) { this.eventStore = eventStore; }
public void setFilterTransactionEntry(boolean filterTransactionEntry) { this.filterTransactionEntry = filterTransactionEntry; }
public void setFilterEmtryTransactionEntry(boolean filterEmtryTransactionEntry) { this.filterEmtryTransactionEntry = filterEmtryTransactionEntry; }
public void setEmptyTransactionInterval(long emptyTransactionInterval) { this.emptyTransactionInterval = emptyTransactionInterval; }
public void setEmptyTransctionThresold(long emptyTransctionThresold) { this.emptyTransctionThresold = emptyTransctionThresold; }
public AtomicLong getEventsSinkBlockingTime() { return eventsSinkBlockingTime; }
}
|