Package org.apache.pulsar.io.canal
Class CanalStringSource
java.lang.Object
org.apache.pulsar.io.core.AbstractPushSource<T>
org.apache.pulsar.io.core.PushSource<V>
org.apache.pulsar.io.canal.CanalAbstractSource<org.apache.pulsar.io.canal.CanalMessage>
org.apache.pulsar.io.canal.CanalStringSource
- All Implemented Interfaces:
AutoCloseable,Source<org.apache.pulsar.io.canal.CanalMessage>
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 TypeMethodDescriptionorg.apache.pulsar.io.canal.CanalMessageextractValue(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
-
CanalStringSource
public CanalStringSource()
-
-
Method Details
-
getMessageId
- Specified by:
getMessageIdin classCanalAbstractSource<org.apache.pulsar.io.canal.CanalMessage>
-
extractValue
public org.apache.pulsar.io.canal.CanalMessage extractValue(List<com.alibaba.otter.canal.protocol.FlatMessage> flatMessages) - Specified by:
extractValuein classCanalAbstractSource<org.apache.pulsar.io.canal.CanalMessage>
-