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>

public class CanalStringSource extends CanalAbstractSource<org.apache.pulsar.io.canal.CanalMessage>
A Simple class for mysql binlog sync to pulsar.
  • Constructor Details

    • CanalStringSource

      public CanalStringSource()
  • Method Details

    • getMessageId

      public Long getMessageId(com.alibaba.otter.canal.protocol.Message message)
      Specified by:
      getMessageId in class CanalAbstractSource<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:
      extractValue in class CanalAbstractSource<org.apache.pulsar.io.canal.CanalMessage>