Package org.apache.pulsar.io.canal
Class CanalByteSource
java.lang.Object
org.apache.pulsar.io.core.AbstractPushSource<T>
org.apache.pulsar.io.core.PushSource<V>
org.apache.pulsar.io.canal.CanalAbstractSource<byte[]>
org.apache.pulsar.io.canal.CanalByteSource
- All Implemented Interfaces:
AutoCloseable,Source<byte[]>
A Simple class for mysql binlog sync to pulsar.
-
Field Summary
Fields inherited from class org.apache.pulsar.io.canal.CanalAbstractSource
handler, running, thread -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionbyte[]extractValue(List<com.alibaba.otter.canal.protocol.FlatMessage> flatMessages) getMessageId(com.alibaba.otter.canal.protocol.Message message) Methods inherited from class org.apache.pulsar.io.canal.CanalAbstractSource
close, open, process, startMethods inherited from class org.apache.pulsar.io.core.PushSource
readMethods inherited from class org.apache.pulsar.io.core.AbstractPushSource
consume, getQueueLength, notifyError, readNext
-
Constructor Details
-
CanalByteSource
public CanalByteSource()
-
-
Method Details
-
getMessageId
- Specified by:
getMessageIdin classCanalAbstractSource<byte[]>
-
extractValue
- Specified by:
extractValuein classCanalAbstractSource<byte[]>
-