diff --git a/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java b/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java new file mode 100644 index 00000000..9d0562b0 --- /dev/null +++ b/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java @@ -0,0 +1,243 @@ +package com.alibaba.otter.canal.instance.core; + +import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; +import com.alibaba.otter.canal.common.alarm.CanalAlarmHandler; +import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter; +import com.alibaba.otter.canal.meta.CanalMetaManager; +import com.alibaba.otter.canal.parse.CanalEventParser; +import com.alibaba.otter.canal.parse.ha.CanalHAController; +import com.alibaba.otter.canal.parse.ha.HeartBeatHAController; +import com.alibaba.otter.canal.parse.inbound.AbstractEventParser; +import com.alibaba.otter.canal.parse.inbound.group.GroupEventParser; +import com.alibaba.otter.canal.parse.inbound.mysql.MysqlEventParser; +import com.alibaba.otter.canal.parse.index.CanalLogPositionManager; +import com.alibaba.otter.canal.protocol.CanalEntry; +import com.alibaba.otter.canal.protocol.ClientIdentity; +import com.alibaba.otter.canal.sink.CanalEventSink; +import com.alibaba.otter.canal.store.CanalEventStore; +import com.alibaba.otter.canal.store.model.Event; +import org.apache.commons.lang.StringUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; + +/** + * Created with Intellig IDEA. + * Author: yinxiu + * Date: 2016-01-07 + * Time: 22:26 + */ +public class AbstractCanalInstance extends AbstractCanalLifeCycle implements CanalInstance { + + private static final Logger logger = LoggerFactory.getLogger(AbstractCanalInstance.class); + + protected Long canalId; // 和manager交互唯一标示 + protected String destination; // 队列名字 + protected CanalEventStore eventStore; // 有序队列 + + protected CanalEventParser eventParser; // 解析对应的数据信息 + protected CanalEventSink> eventSink; // 链接parse和store的桥接器 + protected CanalMetaManager metaManager; // 消费信息管理器 + protected CanalAlarmHandler alarmHandler; // alarm报警机制 + + @Override + public boolean subscribeChange(ClientIdentity identity) { + if (StringUtils.isNotEmpty(identity.getFilter())) { + logger.info("subscribe filter change to " + identity.getFilter()); + AviaterRegexFilter aviaterFilter = new AviaterRegexFilter(identity.getFilter()); + + boolean isGroup = (eventParser instanceof GroupEventParser); + if (isGroup) { + // 处理group的模式 + List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); + for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 + ((AbstractEventParser) singleEventParser).setEventFilter(aviaterFilter); + } + } else { + ((AbstractEventParser) eventParser).setEventFilter(aviaterFilter); + } + + } + + // filter的处理规则 + // a. parser处理数据过滤处理 + // b. sink处理数据的路由&分发,一份parse数据经过sink后可以分发为多份,每份的数据可以根据自己的过滤规则不同而有不同的数据 + // 后续内存版的一对多分发,可以考虑 + return true; + } + + @Override + public void start() { + super.start(); + if (!metaManager.isStart()) { + metaManager.start(); + } + + if (!eventStore.isStart()) { + eventStore.start(); + } + + if (!eventSink.isStart()) { + eventSink.start(); + } + + if (!eventParser.isStart()) { + beforeStartEventParser(eventParser); + eventParser.start(); + afterStartEventParser(eventParser); + } + logger.info("start successful...."); + } + + @Override + public void stop() { + super.stop(); + logger.info("stop CannalInstance for {}-{} ", new Object[] { canalId, destination }); + + if (eventParser.isStart()) { + beforeStopEventParser(eventParser); + eventParser.stop(); + afterStopEventParser(eventParser); + } + + if (eventSink.isStart()) { + eventSink.stop(); + } + + if (eventStore.isStart()) { + eventStore.stop(); + } + + if (metaManager.isStart()) { + metaManager.stop(); + } + + if (alarmHandler.isStart()) { + alarmHandler.stop(); + } + + + logger.info("stop successful...."); + } + + protected void beforeStartEventParser(CanalEventParser eventParser) { + + boolean isGroup = (eventParser instanceof GroupEventParser); + if (isGroup) { + // 处理group的模式 + List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); + for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 + startEventParserInternal(singleEventParser, true); + } + } else { + startEventParserInternal(eventParser, false); + } + } + + // around event parser, default impl + protected void afterStartEventParser(CanalEventParser eventParser) { + // 读取一下历史订阅的filter信息 + List clientIdentitys = metaManager.listAllSubscribeInfo(destination); + for (ClientIdentity clientIdentity : clientIdentitys) { + subscribeChange(clientIdentity); + } + } + + // around event parser + protected void beforeStopEventParser(CanalEventParser eventParser) { + // noop + } + + protected void afterStopEventParser(CanalEventParser eventParser) { + + boolean isGroup = (eventParser instanceof GroupEventParser); + if (isGroup) { + // 处理group的模式 + List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); + for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 + stopEventParserInternal(singleEventParser); + } + } else { + stopEventParserInternal(eventParser); + } + } + + /** + * 初始化单个eventParser,不需要考虑group + */ + protected void startEventParserInternal(CanalEventParser eventParser, boolean isGroup) { + if (eventParser instanceof AbstractEventParser) { + AbstractEventParser abstractEventParser = (AbstractEventParser) eventParser; + // 首先启动log position管理器 + CanalLogPositionManager logPositionManager = abstractEventParser.getLogPositionManager(); + if (!logPositionManager.isStart()) { + logPositionManager.start(); + } + } + + if (eventParser instanceof MysqlEventParser) { + MysqlEventParser mysqlEventParser = (MysqlEventParser) eventParser; + CanalHAController haController = mysqlEventParser.getHaController(); + + if (haController instanceof HeartBeatHAController) { + ((HeartBeatHAController) haController).setCanalHASwitchable(mysqlEventParser); + } + + if (!haController.isStart()) { + haController.start(); + } + + } + } + + protected void stopEventParserInternal(CanalEventParser eventParser) { + if (eventParser instanceof AbstractEventParser) { + AbstractEventParser abstractEventParser = (AbstractEventParser) eventParser; + // 首先启动log position管理器 + CanalLogPositionManager logPositionManager = abstractEventParser.getLogPositionManager(); + if (logPositionManager.isStart()) { + logPositionManager.stop(); + } + } + + if (eventParser instanceof MysqlEventParser) { + MysqlEventParser mysqlEventParser = (MysqlEventParser) eventParser; + CanalHAController haController = mysqlEventParser.getHaController(); + if (haController.isStart()) { + haController.stop(); + } + } + } + + // ==================getter================================== + @Override + public String getDestination() { + return destination; + } + + @Override + public CanalEventParser getEventParser() { + return eventParser; + } + + @Override + public CanalEventSink getEventSink() { + return eventSink; + } + + @Override + public CanalEventStore getEventStore() { + return eventStore; + } + + @Override + public CanalMetaManager getMetaManager() { + return metaManager; + } + + @Override + public CanalAlarmHandler getAlarmHandler() { + return alarmHandler; + } +} diff --git a/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/CanalInstanceSupport.java b/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/CanalInstanceSupport.java deleted file mode 100644 index 8284c31d..00000000 --- a/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/CanalInstanceSupport.java +++ /dev/null @@ -1,105 +0,0 @@ -package com.alibaba.otter.canal.instance.core; - -import java.util.List; - -import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; -import com.alibaba.otter.canal.parse.CanalEventParser; -import com.alibaba.otter.canal.parse.ha.CanalHAController; -import com.alibaba.otter.canal.parse.ha.HeartBeatHAController; -import com.alibaba.otter.canal.parse.inbound.AbstractEventParser; -import com.alibaba.otter.canal.parse.inbound.group.GroupEventParser; -import com.alibaba.otter.canal.parse.inbound.mysql.MysqlEventParser; -import com.alibaba.otter.canal.parse.index.CanalLogPositionManager; - -/** - * @author zebin.xuzb 2012-10-17 下午3:12:34 - * @version 1.0.0 - */ -public abstract class CanalInstanceSupport extends AbstractCanalLifeCycle { - - protected void beforeStartEventParser(CanalEventParser eventParser) { - - boolean isGroup = (eventParser instanceof GroupEventParser); - if (isGroup) { - // 处理group的模式 - List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); - for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 - startEventParserInternal(singleEventParser, true); - } - } else { - startEventParserInternal(eventParser, false); - } - } - - // around event parser - protected void afterStartEventParser(CanalEventParser eventParser) { - // noop - } - - // around event parser - protected void beforeStopEventParser(CanalEventParser eventParser) { - // noop - } - - protected void afterStopEventParser(CanalEventParser eventParser) { - - boolean isGroup = (eventParser instanceof GroupEventParser); - if (isGroup) { - // 处理group的模式 - List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); - for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 - stopEventParserInternal(singleEventParser); - } - } else { - stopEventParserInternal(eventParser); - } - } - - /** - * 初始化单个eventParser,不需要考虑group - */ - protected void startEventParserInternal(CanalEventParser eventParser, boolean isGroup) { - if (eventParser instanceof AbstractEventParser) { - AbstractEventParser abstractEventParser = (AbstractEventParser) eventParser; - // 首先启动log position管理器 - CanalLogPositionManager logPositionManager = abstractEventParser.getLogPositionManager(); - if (!logPositionManager.isStart()) { - logPositionManager.start(); - } - } - - if (eventParser instanceof MysqlEventParser) { - MysqlEventParser mysqlEventParser = (MysqlEventParser) eventParser; - CanalHAController haController = mysqlEventParser.getHaController(); - - if (haController instanceof HeartBeatHAController) { - ((HeartBeatHAController) haController).setCanalHASwitchable(mysqlEventParser); - } - - if (!haController.isStart()) { - haController.start(); - } - - } - } - - protected void stopEventParserInternal(CanalEventParser eventParser) { - if (eventParser instanceof AbstractEventParser) { - AbstractEventParser abstractEventParser = (AbstractEventParser) eventParser; - // 首先启动log position管理器 - CanalLogPositionManager logPositionManager = abstractEventParser.getLogPositionManager(); - if (logPositionManager.isStart()) { - logPositionManager.stop(); - } - } - - if (eventParser instanceof MysqlEventParser) { - MysqlEventParser mysqlEventParser = (MysqlEventParser) eventParser; - CanalHAController haController = mysqlEventParser.getHaController(); - if (haController.isStart()) { - haController.stop(); - } - } - } - -} diff --git a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java index 5e979a5e..07bbff83 100644 --- a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java +++ b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java @@ -6,6 +6,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import com.alibaba.otter.canal.instance.core.AbstractCanalInstance; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -18,7 +19,6 @@ import com.alibaba.otter.canal.common.utils.JsonUtils; import com.alibaba.otter.canal.common.zookeeper.ZkClientx; import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter; import com.alibaba.otter.canal.instance.core.CanalInstance; -import com.alibaba.otter.canal.instance.core.CanalInstanceSupport; import com.alibaba.otter.canal.instance.manager.model.Canal; import com.alibaba.otter.canal.instance.manager.model.CanalParameter; import com.alibaba.otter.canal.instance.manager.model.CanalParameter.DataSourcing; @@ -64,10 +64,9 @@ import com.alibaba.otter.canal.store.model.Event; * @author jianghang 2012-7-11 下午09:26:51 * @version 1.0.0 */ -public class CanalInstanceWithManager extends CanalInstanceSupport implements CanalInstance { +public class CanalInstanceWithManager extends AbstractCanalInstance { private static final Logger logger = LoggerFactory.getLogger(CanalInstanceWithManager.class); - protected Long canalId; // 和manager交互唯一标示 protected String destination; // 队列名字 protected String filter; // 过滤表达式 protected CanalParameter parameters; // 对应参数 @@ -77,11 +76,6 @@ public class CanalInstanceWithManager extends CanalInstanceSupport implements Ca protected CanalEventParser eventParser; // 解析对应的数据信息 protected CanalEventSink> eventSink; // 链接parse和store的桥接器 protected CanalAlarmHandler alarmHandler; // alarm报警机制 - protected ZkClientx zkClientx; - - public CanalInstanceWithManager(Canal canal){ - this(canal, null); - } public CanalInstanceWithManager(Canal canal, String filter){ this.parameters = canal.getCanalParameter(); @@ -89,7 +83,7 @@ public class CanalInstanceWithManager extends CanalInstanceSupport implements Ca this.destination = canal.getName(); this.filter = filter; - logger.info("init CannalInstance for {}-{} with parameters:{}", canalId, destination, parameters); + logger.info("init CanalInstance for {}-{} with parameters:{}", canalId, destination, parameters); // 初始化报警机制 initAlarmHandler(); // 初始化metaManager @@ -101,7 +95,7 @@ public class CanalInstanceWithManager extends CanalInstanceSupport implements Ca // 初始化eventParser; initEventParser(); - // 基础工具,需要提前start,会有先订阅再根据filter条件启动paser的需求 + // 基础工具,需要提前start,会有先订阅再根据filter条件启动parse的需求 if (!alarmHandler.isStart()) { alarmHandler.start(); } @@ -113,98 +107,9 @@ public class CanalInstanceWithManager extends CanalInstanceSupport implements Ca } public void start() { - super.start(); // 初始化metaManager logger.info("start CannalInstance for {}-{} with parameters:{}", canalId, destination, parameters); - - if (!metaManager.isStart()) { - metaManager.start(); - } - - if (!alarmHandler.isStart()) { - alarmHandler.start(); - } - - if (!eventStore.isStart()) { - eventStore.start(); - } - - if (!eventSink.isStart()) { - eventSink.start(); - } - - if (!eventParser.isStart()) { - beforeStartEventParser(eventParser); - eventParser.start(); - } - - logger.info("start successful...."); - } - - public void stop() { - logger.info("stop CannalInstance for {}-{} ", new Object[] { canalId, destination }); - - if (eventParser.isStart()) { - eventParser.stop(); - afterStopEventParser(eventParser); - } - - if (eventSink.isStart()) { - eventSink.stop(); - } - - if (eventStore.isStart()) { - eventStore.stop(); - } - - if (metaManager.isStart()) { - metaManager.stop(); - } - - if (alarmHandler.isStart()) { - alarmHandler.stop(); - } - - // if (zkClientx != null) { - // zkClientx.close(); - // } - - super.stop(); - logger.info("stop successful...."); - } - - public boolean subscribeChange(ClientIdentity identity) { - if (StringUtils.isNotEmpty(identity.getFilter())) { - AviaterRegexFilter aviaterFilter = new AviaterRegexFilter(identity.getFilter()); - - boolean isGroup = (eventParser instanceof GroupEventParser); - if (isGroup) { - // 处理group的模式 - List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); - for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 - ((AbstractEventParser) singleEventParser).setEventFilter(aviaterFilter); - } - } else { - ((AbstractEventParser) eventParser).setEventFilter(aviaterFilter); - } - - } - - // filter的处理规则 - // a. parser处理数据过滤处理 - // b. sink处理数据的路由&分发,一份parse数据经过sink后可以分发为多份,每份的数据可以根据自己的过滤规则不同而有不同的数据 - // 后续内存版的一对多分发,可以考虑 - return true; - } - - protected void afterStartEventParser(CanalEventParser eventParser) { - super.afterStartEventParser(eventParser); - - // 读取一下历史订阅的filter信息 - List clientIdentitys = metaManager.listAllSubscribeInfo(destination); - for (ClientIdentity clientIdentity : clientIdentitys) { - subscribeChange(clientIdentity); - } + super.start(); } protected void initAlarmHandler() { @@ -515,34 +420,4 @@ public class CanalInstanceWithManager extends CanalInstanceSupport implements Ca return ZkClientx.getZkClient(StringUtils.join(zkClusters, ";")); } - // ===================================== - - public String getDestination() { - return destination; - } - - public CanalMetaManager getMetaManager() { - return metaManager; - } - - public CanalEventStore getEventStore() { - return eventStore; - } - - public CanalEventParser getEventParser() { - return eventParser; - } - - public CanalEventSink> getEventSink() { - return eventSink; - } - - public CanalAlarmHandler getAlarmHandler() { - return alarmHandler; - } - - public void setAlarmHandler(CanalAlarmHandler alarmHandler) { - this.alarmHandler = alarmHandler; - } - } diff --git a/instance/spring/src/main/java/com/alibaba/otter/canal/instance/spring/CanalInstanceWithSpring.java b/instance/spring/src/main/java/com/alibaba/otter/canal/instance/spring/CanalInstanceWithSpring.java index 20b553d8..81ad2ff0 100644 --- a/instance/spring/src/main/java/com/alibaba/otter/canal/instance/spring/CanalInstanceWithSpring.java +++ b/instance/spring/src/main/java/com/alibaba/otter/canal/instance/spring/CanalInstanceWithSpring.java @@ -2,20 +2,14 @@ package com.alibaba.otter.canal.instance.spring; import java.util.List; -import org.apache.commons.lang.StringUtils; +import com.alibaba.otter.canal.instance.core.AbstractCanalInstance; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.alibaba.otter.canal.common.alarm.CanalAlarmHandler; -import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter; -import com.alibaba.otter.canal.instance.core.CanalInstance; -import com.alibaba.otter.canal.instance.core.CanalInstanceSupport; import com.alibaba.otter.canal.meta.CanalMetaManager; import com.alibaba.otter.canal.parse.CanalEventParser; -import com.alibaba.otter.canal.parse.inbound.AbstractEventParser; -import com.alibaba.otter.canal.parse.inbound.group.GroupEventParser; import com.alibaba.otter.canal.protocol.CanalEntry; -import com.alibaba.otter.canal.protocol.ClientIdentity; import com.alibaba.otter.canal.sink.CanalEventSink; import com.alibaba.otter.canal.store.CanalEventStore; import com.alibaba.otter.canal.store.model.Event; @@ -27,130 +21,13 @@ import com.alibaba.otter.canal.store.model.Event; * @author zebin.xuzb * @version 1.0.0 */ -public class CanalInstanceWithSpring extends CanalInstanceSupport implements CanalInstance { +public class CanalInstanceWithSpring extends AbstractCanalInstance { private static final Logger logger = LoggerFactory.getLogger(CanalInstanceWithSpring.class); - private String destination; - private CanalEventParser eventParser; - private CanalEventSink> eventSink; - private CanalEventStore eventStore; - private CanalMetaManager metaManager; - private CanalAlarmHandler alarmHandler; - - public String getDestination() { - return this.destination; - } - - public CanalEventParser getEventParser() { - return this.eventParser; - } - - public CanalEventSink> getEventSink() { - return this.eventSink; - } - - public CanalEventStore getEventStore() { - return this.eventStore; - } - - public CanalMetaManager getMetaManager() { - return this.metaManager; - } - - public CanalAlarmHandler getAlarmHandler() { - return alarmHandler; - } - - public boolean subscribeChange(ClientIdentity identity) { - if (StringUtils.isNotEmpty(identity.getFilter())) { - logger.info("subscribe filter change to " + identity.getFilter()); - AviaterRegexFilter aviaterFilter = new AviaterRegexFilter(identity.getFilter()); - - boolean isGroup = (eventParser instanceof GroupEventParser); - if (isGroup) { - // 处理group的模式 - List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); - for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 - ((AbstractEventParser) singleEventParser).setEventFilter(aviaterFilter); - } - } else { - ((AbstractEventParser) eventParser).setEventFilter(aviaterFilter); - } - - } - - // filter的处理规则 - // a. parser处理数据过滤处理 - // b. sink处理数据的路由&分发,一份parse数据经过sink后可以分发为多份,每份的数据可以根据自己的过滤规则不同而有不同的数据 - // 后续内存版的一对多分发,可以考虑 - return true; - } - - protected void afterStartEventParser(CanalEventParser eventParser) { - super.afterStartEventParser(eventParser); - - // 读取一下历史订阅的filter信息 - List clientIdentitys = metaManager.listAllSubscribeInfo(destination); - for (ClientIdentity clientIdentity : clientIdentitys) { - subscribeChange(clientIdentity); - } - } public void start() { - super.start(); - logger.info("start CannalInstance for {}-{} ", new Object[] { 1, destination }); - if (!metaManager.isStart()) { - metaManager.start(); - } - - if (!eventStore.isStart()) { - eventStore.start(); - } - - if (!eventSink.isStart()) { - eventSink.start(); - } - - if (!eventParser.isStart()) { - beforeStartEventParser(eventParser); - eventParser.start(); - afterStartEventParser(eventParser); - } - - logger.info("start successful...."); - } - - public void stop() { - logger.info("stop CannalInstance for {}-{} ", new Object[] { 1, destination }); - if (eventParser.isStart()) { - beforeStopEventParser(eventParser); - eventParser.stop(); - afterStopEventParser(eventParser); - } - - if (eventSink.isStart()) { - eventSink.stop(); - } - - if (eventStore.isStart()) { - eventStore.stop(); - } - - if (metaManager.isStart()) { - metaManager.stop(); - } - - if (alarmHandler.isStart()) { - alarmHandler.stop(); - } - - // if (zkClientx != null) { - // zkClientx.close(); - // } - - super.stop(); - logger.info("stop successful...."); + super.start(); } // ======== setter ======== diff --git a/server/src/main/java/com/alibaba/otter/canal/server/netty/CanalServerWithNetty.java b/server/src/main/java/com/alibaba/otter/canal/server/netty/CanalServerWithNetty.java index 3df5a15f..8876e6b1 100644 --- a/server/src/main/java/com/alibaba/otter/canal/server/netty/CanalServerWithNetty.java +++ b/server/src/main/java/com/alibaba/otter/canal/server/netty/CanalServerWithNetty.java @@ -3,6 +3,7 @@ package com.alibaba.otter.canal.server.netty; import java.net.InetSocketAddress; import java.util.concurrent.Executors; +import com.alibaba.otter.canal.instance.core.CanalInstance; import org.apache.commons.lang.StringUtils; import org.jboss.netty.bootstrap.ServerBootstrap; import org.jboss.netty.channel.Channel;