Class CanalAbstractSource<V>

java.lang.Object
org.apache.pulsar.io.core.PushSource<V>
org.apache.pulsar.io.canal.CanalAbstractSource<V>
All Implemented Interfaces:
AutoCloseable, org.apache.pulsar.io.core.Source<V>
Direct Known Subclasses:
CanalByteSource, CanalStringSource

public abstract class CanalAbstractSource<V> extends org.apache.pulsar.io.core.PushSource<V>
A Simple abstract class for mysql binlog sync to pulsar.
  • Field Details

  • Constructor Details

    • CanalAbstractSource

      public CanalAbstractSource()
  • Method Details

    • open

      public void open(Map<String,Object> config, org.apache.pulsar.io.core.SourceContext sourceContext) throws Exception
      Specified by:
      open in interface org.apache.pulsar.io.core.Source<V>
      Specified by:
      open in class org.apache.pulsar.io.core.PushSource<V>
      Throws:
      Exception
    • start

      protected void start()
    • close

      public void close() throws InterruptedException
      Throws:
      InterruptedException
    • process

      protected void process()
    • getMessageId

      public abstract Long getMessageId(com.alibaba.otter.canal.protocol.Message message)
    • extractValue

      public abstract V extractValue(List<com.alibaba.otter.canal.protocol.FlatMessage> flatMessages)