diff --git a/common/src/main/java/com/alibaba/otter/canal/common/zookeeper/ZkClientx.java b/common/src/main/java/com/alibaba/otter/canal/common/zookeeper/ZkClientx.java index cba75cc3..74bde300 100644 --- a/common/src/main/java/com/alibaba/otter/canal/common/zookeeper/ZkClientx.java +++ b/common/src/main/java/com/alibaba/otter/canal/common/zookeeper/ZkClientx.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.common.zookeeper; +import com.google.common.collect.MigrateMap; import java.util.Map; import org.I0Itec.zkclient.IZkConnection; @@ -23,12 +24,14 @@ import com.google.common.collect.MapMaker; public class ZkClientx extends ZkClient { // 对于zkclient进行一次缓存,避免一个jvm内部使用多个zk connection - private static Map clients = new MapMaker().makeComputingMap(new Function() { + private static Map clients = MigrateMap.makeComputingMap(new Function() + { - public ZkClientx apply(String servers) { - return new ZkClientx(servers); - } - }); + public ZkClientx apply(String servers) + { + return new ZkClientx(servers); + } + }); public static ZkClientx getZkClient(String servers) { return clients.get(servers); diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalController.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalController.java index 3e8f0c8a..8ba01a98 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalController.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalController.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.deployer; +import com.google.common.collect.MigrateMap; import java.util.Map; import java.util.Properties; @@ -70,9 +71,11 @@ public class CanalController { } public CanalController(final Properties properties){ - managerClients = new MapMaker().makeComputingMap(new Function() { + managerClients = MigrateMap.makeComputingMap(new Function() + { - public CanalConfigClient apply(String managerAddress) { + public CanalConfigClient apply(String managerAddress) + { return getManagerClient(managerAddress); } }); @@ -107,68 +110,88 @@ public class CanalController { final ServerRunningData serverData = new ServerRunningData(cid, ip + ":" + port); ServerRunningMonitors.setServerData(serverData); - ServerRunningMonitors.setRunningMonitors(new MapMaker().makeComputingMap(new Function() { - - public ServerRunningMonitor apply(final String destination) { + ServerRunningMonitors.setRunningMonitors(MigrateMap.makeComputingMap(new Function() + { + public ServerRunningMonitor apply(final String destination) + { ServerRunningMonitor runningMonitor = new ServerRunningMonitor(serverData); runningMonitor.setDestination(destination); - runningMonitor.setListener(new ServerRunningListener() { + runningMonitor.setListener(new ServerRunningListener() + { - public void processActiveEnter() { - try { + public void processActiveEnter() + { + try + { MDC.put(CanalConstants.MDC_DESTINATION, String.valueOf(destination)); embededCanalServer.start(destination); - } finally { + } finally + { MDC.remove(CanalConstants.MDC_DESTINATION); } } - public void processActiveExit() { - try { + public void processActiveExit() + { + try + { MDC.put(CanalConstants.MDC_DESTINATION, String.valueOf(destination)); embededCanalServer.stop(destination); - } finally { + } finally + { MDC.remove(CanalConstants.MDC_DESTINATION); } } - public void processStart() { - try { - if (zkclientx != null) { + public void processStart() + { + try + { + if (zkclientx != null) + { final String path = ZookeeperPathUtils.getDestinationClusterNode(destination, ip + ":" - + port); + + port); initCid(path); - zkclientx.subscribeStateChanges(new IZkStateListener() { + zkclientx.subscribeStateChanges(new IZkStateListener() + { - public void handleStateChanged(KeeperState state) throws Exception { + public void handleStateChanged(KeeperState state) throws Exception + { } - public void handleNewSession() throws Exception { + public void handleNewSession() throws Exception + { initCid(path); } }); } - } finally { + } finally + { MDC.remove(CanalConstants.MDC_DESTINATION); } } - public void processStop() { - try { + public void processStop() + { + try + { MDC.put(CanalConstants.MDC_DESTINATION, String.valueOf(destination)); - if (zkclientx != null) { + if (zkclientx != null) + { final String path = ZookeeperPathUtils.getDestinationClusterNode(destination, ip + ":" - + port); + + port); releaseCid(path); } - } finally { + } finally + { MDC.remove(CanalConstants.MDC_DESTINATION); } } }); - if (zkclientx != null) { + if (zkclientx != null) + { runningMonitor.setZkClient(zkclientx); } return runningMonitor; @@ -216,25 +239,32 @@ public class CanalController { } }; - instanceConfigMonitors = new MapMaker().makeComputingMap(new Function() { + instanceConfigMonitors = MigrateMap.makeComputingMap(new Function() + { - public InstanceConfigMonitor apply(InstanceMode mode) { - int scanInterval = Integer.valueOf(getProperty(properties, CanalConstants.CANAL_AUTO_SCAN_INTERVAL)); + public InstanceConfigMonitor apply(InstanceMode mode) + { + int scanInterval = Integer + .valueOf(getProperty(properties, CanalConstants.CANAL_AUTO_SCAN_INTERVAL)); - if (mode.isSpring()) { + if (mode.isSpring()) + { SpringInstanceConfigMonitor monitor = new SpringInstanceConfigMonitor(); monitor.setScanIntervalInSecond(scanInterval); monitor.setDefaultAction(defaultAction); // 设置conf目录,默认是user.dir + conf目录组成 String rootDir = getProperty(properties, CanalConstants.CANAL_CONF_DIR); - if (StringUtils.isEmpty(rootDir)) { + if (StringUtils.isEmpty(rootDir)) + { rootDir = "../conf"; } monitor.setRootConf(rootDir); return monitor; - } else if (mode.isManager()) { + } else if (mode.isManager()) + { return new ManagerInstanceConfigMonitor(); - } else { + } else + { throw new UnsupportedOperationException("unknow mode :" + mode + " for monitor"); } } diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/SpringInstanceConfigMonitor.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/SpringInstanceConfigMonitor.java index c56edf21..d3b769e1 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/SpringInstanceConfigMonitor.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/monitor/SpringInstanceConfigMonitor.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.deployer.monitor; +import com.google.common.collect.MigrateMap; import java.io.File; import java.io.FileFilter; import java.io.FilenameFilter; @@ -39,12 +40,14 @@ public class SpringInstanceConfigMonitor extends AbstractCanalLifeCycle implemen private long scanIntervalInSecond = 5; private InstanceAction defaultAction = null; private Map actions = new MapMaker().makeMap(); - private Map lastFiles = new MapMaker().makeComputingMap(new Function() { + private Map lastFiles = MigrateMap.makeComputingMap(new Function() + { - public InstanceConfigFiles apply(String destination) { - return new InstanceConfigFiles(destination); - } - }); + public InstanceConfigFiles apply(String destination) + { + return new InstanceConfigFiles(destination); + } + }); private ScheduledExecutorService executor = Executors.newScheduledThreadPool(1, new NamedThreadFactory("canal-instance-scan")); diff --git a/filter/src/main/java/com/alibaba/otter/canal/filter/PatternUtils.java b/filter/src/main/java/com/alibaba/otter/canal/filter/PatternUtils.java index 4ce36315..ce10dc3f 100644 --- a/filter/src/main/java/com/alibaba/otter/canal/filter/PatternUtils.java +++ b/filter/src/main/java/com/alibaba/otter/canal/filter/PatternUtils.java @@ -1,48 +1,54 @@ package com.alibaba.otter.canal.filter; +import com.alibaba.otter.canal.filter.exception.CanalFilterException; +import com.google.common.base.Function; +import com.google.common.collect.MapMaker; +import com.google.common.collect.MigrateMap; import java.util.Map; - import org.apache.oro.text.regex.MalformedPatternException; import org.apache.oro.text.regex.Pattern; import org.apache.oro.text.regex.PatternCompiler; import org.apache.oro.text.regex.Perl5Compiler; -import com.alibaba.otter.canal.filter.exception.CanalFilterException; -import com.google.common.base.Function; -import com.google.common.collect.MapMaker; - /** * 提供{@linkplain Pattern}的lazy get处理 - * + * * @author jianghang 2013-1-22 下午09:36:44 * @version 1.0.0 */ -public class PatternUtils { +public class PatternUtils +{ - private static Map patterns = new MapMaker().softValues().makeComputingMap( - new Function() { + private static Map patterns = MigrateMap.makeComputingMap(new MapMaker().softValues(), + new Function() + { - public Pattern apply( - String pattern) { - try { - PatternCompiler pc = new Perl5Compiler(); - return pc.compile( - pattern, - Perl5Compiler.CASE_INSENSITIVE_MASK - | Perl5Compiler.READ_ONLY_MASK - | Perl5Compiler.SINGLELINE_MASK); - } catch (MalformedPatternException e) { - throw new CanalFilterException( - e); - } - } - }); + public Pattern apply( + String pattern) + { + try + { + PatternCompiler pc = new Perl5Compiler(); + return pc.compile( + pattern, + Perl5Compiler.CASE_INSENSITIVE_MASK + | Perl5Compiler.READ_ONLY_MASK + | Perl5Compiler.SINGLELINE_MASK); + } catch (MalformedPatternException e) + { + throw new CanalFilterException( + e); + } + } + }); - public static Pattern getPattern(String pattern) { + public static Pattern getPattern(String pattern) + { return patterns.get(pattern); } - public static void clear() { + public static void clear() + { patterns.clear(); } } diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java index f7062c53..7319c2cb 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/FileMixedMetaManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.meta; +import com.google.common.collect.MigrateMap; import java.io.File; import java.io.IOException; import java.nio.charset.Charset; @@ -69,28 +70,36 @@ public class FileMixedMetaManager extends MemoryMetaManager implements CanalMeta throw new CanalMetaManagerException("dir[" + dataDir.getPath() + "] can not read/write"); } - dataFileCaches = new MapMaker().makeComputingMap(new Function() { + dataFileCaches = MigrateMap.makeComputingMap(new Function() + { - public File apply(String destination) { + public File apply(String destination) + { return getDataFile(destination); } }); executor = Executors.newScheduledThreadPool(1); - destinations = new MapMaker().makeComputingMap(new Function>() { + destinations = MigrateMap.makeComputingMap(new Function>() + { - public List apply(String destination) { + public List apply(String destination) + { return loadClientIdentity(destination); } }); - cursors = new MapMaker().makeComputingMap(new Function() { + cursors = MigrateMap.makeComputingMap(new Function() + { - public Position apply(ClientIdentity clientIdentity) { + public Position apply(ClientIdentity clientIdentity) + { Position position = loadCursor(clientIdentity.getDestination(), clientIdentity); - if (position == null) { + if (position == null) + { return nullCursor; // 返回一个空对象标识,避免出现异常 - } else { + } else + { return position; } } diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java index 6effab1a..f7cbb6b4 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/MemoryMetaManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.meta; +import com.google.common.collect.MigrateMap; import java.util.Collections; import java.util.List; import java.util.Map; @@ -31,9 +32,11 @@ public class MemoryMetaManager extends AbstractCanalLifeCycle implements CanalMe public void start() { super.start(); - batches = new MapMaker().makeComputingMap(new Function() { + batches = MigrateMap.makeComputingMap(new Function() + { - public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) { + public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) + { return MemoryClientIdentityBatch.create(clientIdentity); } @@ -41,9 +44,11 @@ public class MemoryMetaManager extends AbstractCanalLifeCycle implements CanalMe cursors = new MapMaker().makeMap(); - destinations = new MapMaker().makeComputingMap(new Function>() { + destinations = MigrateMap.makeComputingMap(new Function>() + { - public List apply(String destination) { + public List apply(String destination) + { return Lists.newArrayList(); } }); diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/MixedMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/MixedMetaManager.java index 84d7696d..f835a156 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/MixedMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/MixedMetaManager.java @@ -1,68 +1,79 @@ package com.alibaba.otter.canal.meta; -import java.util.List; -import java.util.Map; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.springframework.util.Assert; - import com.alibaba.otter.canal.meta.exception.CanalMetaManagerException; import com.alibaba.otter.canal.protocol.ClientIdentity; import com.alibaba.otter.canal.protocol.position.Position; import com.alibaba.otter.canal.protocol.position.PositionRange; import com.google.common.base.Function; -import com.google.common.collect.MapMaker; +import com.google.common.collect.MigrateMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import org.springframework.util.Assert; /** * 组合memory + zookeeper的使用模式 - * + * * @author jianghang 2012-7-11 下午03:58:00 * @version 1.0.0 */ -public class MixedMetaManager extends MemoryMetaManager implements CanalMetaManager { +public class MixedMetaManager extends MemoryMetaManager implements CanalMetaManager +{ - private ExecutorService executor; + private ExecutorService executor; private ZooKeeperMetaManager zooKeeperMetaManager; @SuppressWarnings("serial") - private final Position nullCursor = new Position() { - }; + private final Position nullCursor = new Position() + { + }; - public void start() { + public void start() + { super.start(); Assert.notNull(zooKeeperMetaManager); - if (!zooKeeperMetaManager.isStart()) { + if (!zooKeeperMetaManager.isStart()) + { zooKeeperMetaManager.start(); } executor = Executors.newFixedThreadPool(1); - destinations = new MapMaker().makeComputingMap(new Function>() { + destinations = MigrateMap.makeComputingMap(new Function>() + { - public List apply(String destination) { + public List apply(String destination) + { return zooKeeperMetaManager.listAllSubscribeInfo(destination); } }); - cursors = new MapMaker().makeComputingMap(new Function() { + cursors = MigrateMap.makeComputingMap(new Function() + { - public Position apply(ClientIdentity clientIdentity) { + public Position apply(ClientIdentity clientIdentity) + { Position position = zooKeeperMetaManager.getCursor(clientIdentity); - if (position == null) { + if (position == null) + { return nullCursor; // 返回一个空对象标识,避免出现异常 - } else { + } else + { return position; } } }); - batches = new MapMaker().makeComputingMap(new Function() { + batches = MigrateMap.makeComputingMap(new Function() + { - public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) { + public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) + { // 读取一下zookeeper信息,初始化一次 MemoryClientIdentityBatch batches = MemoryClientIdentityBatch.create(clientIdentity); Map positionRanges = zooKeeperMetaManager.listAllBatchs(clientIdentity); - for (Map.Entry entry : positionRanges.entrySet()) { + for (Map.Entry entry : positionRanges.entrySet()) + { batches.addPositionRange(entry.getValue(), entry.getKey()); // 添加记录到指定batchId } return batches; @@ -70,10 +81,12 @@ public class MixedMetaManager extends MemoryMetaManager implements CanalMetaMana }); } - public void stop() { + public void stop() + { super.stop(); - if (zooKeeperMetaManager.isStart()) { + if (zooKeeperMetaManager.isStart()) + { zooKeeperMetaManager.stop(); } @@ -82,58 +95,73 @@ public class MixedMetaManager extends MemoryMetaManager implements CanalMetaMana batches.clear(); } - public void subscribe(final ClientIdentity clientIdentity) throws CanalMetaManagerException { + public void subscribe(final ClientIdentity clientIdentity) throws CanalMetaManagerException + { super.subscribe(clientIdentity); - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.subscribe(clientIdentity); } }); } - public void unsubscribe(final ClientIdentity clientIdentity) throws CanalMetaManagerException { + public void unsubscribe(final ClientIdentity clientIdentity) throws CanalMetaManagerException + { super.unsubscribe(clientIdentity); - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.unsubscribe(clientIdentity); } }); } public void updateCursor(final ClientIdentity clientIdentity, final Position position) - throws CanalMetaManagerException { + throws CanalMetaManagerException + { super.updateCursor(clientIdentity, position); // 异步刷新 - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.updateCursor(clientIdentity, position); } }); } @Override - public Position getCursor(ClientIdentity clientIdentity) throws CanalMetaManagerException { + public Position getCursor(ClientIdentity clientIdentity) throws CanalMetaManagerException + { Position position = super.getCursor(clientIdentity); - if (position == nullCursor) { + if (position == nullCursor) + { return null; - } else { + } else + { return position; } } public Long addBatch(final ClientIdentity clientIdentity, final PositionRange positionRange) - throws CanalMetaManagerException { + throws CanalMetaManagerException + { final Long batchId = super.addBatch(clientIdentity, positionRange); // 异步刷新 - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.addBatch(clientIdentity, positionRange, batchId); } }); @@ -141,24 +169,30 @@ public class MixedMetaManager extends MemoryMetaManager implements CanalMetaMana } public void addBatch(final ClientIdentity clientIdentity, final PositionRange positionRange, final Long batchId) - throws CanalMetaManagerException { + throws CanalMetaManagerException + { super.addBatch(clientIdentity, positionRange, batchId); // 异步刷新 - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.addBatch(clientIdentity, positionRange, batchId); } }); } public PositionRange removeBatch(final ClientIdentity clientIdentity, final Long batchId) - throws CanalMetaManagerException { + throws CanalMetaManagerException + { PositionRange positionRange = super.removeBatch(clientIdentity, batchId); // 异步刷新 - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.removeBatch(clientIdentity, batchId); } }); @@ -166,20 +200,24 @@ public class MixedMetaManager extends MemoryMetaManager implements CanalMetaMana return positionRange; } - public void clearAllBatchs(final ClientIdentity clientIdentity) throws CanalMetaManagerException { + public void clearAllBatchs(final ClientIdentity clientIdentity) throws CanalMetaManagerException + { super.clearAllBatchs(clientIdentity); // 异步刷新 - executor.submit(new Runnable() { + executor.submit(new Runnable() + { - public void run() { + public void run() + { zooKeeperMetaManager.clearAllBatchs(clientIdentity); } }); } // =============== setter / getter ================ - public void setZooKeeperMetaManager(ZooKeeperMetaManager zooKeeperMetaManager) { + public void setZooKeeperMetaManager(ZooKeeperMetaManager zooKeeperMetaManager) + { this.zooKeeperMetaManager = zooKeeperMetaManager; } } diff --git a/meta/src/main/java/com/alibaba/otter/canal/meta/PeriodMixedMetaManager.java b/meta/src/main/java/com/alibaba/otter/canal/meta/PeriodMixedMetaManager.java index c4c978d7..24fbce0d 100644 --- a/meta/src/main/java/com/alibaba/otter/canal/meta/PeriodMixedMetaManager.java +++ b/meta/src/main/java/com/alibaba/otter/canal/meta/PeriodMixedMetaManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.meta; +import com.google.common.collect.MigrateMap; import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; @@ -52,32 +53,40 @@ public class PeriodMixedMetaManager extends MemoryMetaManager implements CanalMe } executor = Executors.newScheduledThreadPool(1); - destinations = new MapMaker().makeComputingMap(new Function>() { + destinations = MigrateMap.makeComputingMap(new Function>() + { - public List apply(String destination) { + public List apply(String destination) + { return zooKeeperMetaManager.listAllSubscribeInfo(destination); } }); - cursors = new MapMaker().makeComputingMap(new Function() { + cursors = MigrateMap.makeComputingMap(new Function() + { - public Position apply(ClientIdentity clientIdentity) { + public Position apply(ClientIdentity clientIdentity) + { Position position = zooKeeperMetaManager.getCursor(clientIdentity); - if (position == null) { + if (position == null) + { return nullCursor; // 返回一个空对象标识,避免出现异常 - } else { + } else + { return position; } } }); - batches = new MapMaker().makeComputingMap(new Function() { - - public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) { + batches = MigrateMap.makeComputingMap(new Function() + { + public MemoryClientIdentityBatch apply(ClientIdentity clientIdentity) + { // 读取一下zookeeper信息,初始化一次 MemoryClientIdentityBatch batches = MemoryClientIdentityBatch.create(clientIdentity); Map positionRanges = zooKeeperMetaManager.listAllBatchs(clientIdentity); - for (Map.Entry entry : positionRanges.entrySet()) { + for (Map.Entry entry : positionRanges.entrySet()) + { batches.addPositionRange(entry.getValue(), entry.getKey()); // 添加记录到指定batchId } return batches; diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java index f731ff0c..c0f34336 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.parse.inbound.mysql.dbsync; +import com.google.common.collect.MigrateMap; import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; @@ -38,17 +39,23 @@ public class TableMetaCache { public TableMetaCache(MysqlConnection con){ this.connection = con; - tableMetaCache = new MapMaker().makeComputingMap(new Function() { + tableMetaCache = MigrateMap.makeComputingMap(new Function() + { - public TableMeta apply(String name) { - try { + public TableMeta apply(String name) + { + try + { return getTableMeta0(name); - } catch (IOException e) { + } catch (IOException e) + { // 尝试做一次retry操作 - try { + try + { connection.reconnect(); return getTableMeta0(name); - } catch (IOException e1) { + } catch (IOException e1) + { throw new CanalParseException("fetch failed by table meta:" + name, e1); } } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/index/FileMixedLogPositionManager.java b/parse/src/main/java/com/alibaba/otter/canal/parse/index/FileMixedLogPositionManager.java index ee1e178b..d5696857 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/index/FileMixedLogPositionManager.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/index/FileMixedLogPositionManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.parse.index; +import com.google.common.collect.MigrateMap; import java.io.File; import java.io.IOException; import java.nio.charset.Charset; @@ -67,21 +68,27 @@ public class FileMixedLogPositionManager extends MemoryLogPositionManager { throw new CanalMetaManagerException("dir[" + dataDir.getPath() + "] can not read/write"); } - dataFileCaches = new MapMaker().makeComputingMap(new Function() { + dataFileCaches = MigrateMap.makeComputingMap(new Function() + { - public File apply(String destination) { + public File apply(String destination) + { return getDataFile(destination); } }); executor = Executors.newScheduledThreadPool(1); - positions = new MapMaker().makeComputingMap(new Function() { + positions = MigrateMap.makeComputingMap(new Function() + { - public LogPosition apply(String destination) { + public LogPosition apply(String destination) + { LogPosition logPosition = loadDataFromFile(dataFileCaches.get(destination)); - if (logPosition == null) { + if (logPosition == null) + { return nullPosition; - } else { + } else + { return logPosition; } } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/index/MixedLogPositionManager.java b/parse/src/main/java/com/alibaba/otter/canal/parse/index/MixedLogPositionManager.java index 6dfc7c76..161da8e2 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/index/MixedLogPositionManager.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/index/MixedLogPositionManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.parse.index; +import com.google.common.collect.MigrateMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -34,13 +35,17 @@ public class MixedLogPositionManager extends MemoryLogPositionManager implements zooKeeperLogPositionManager.start(); } executor = Executors.newFixedThreadPool(1); - positions = new MapMaker().makeComputingMap(new Function() { + positions = MigrateMap.makeComputingMap(new Function() + { - public LogPosition apply(String destination) { + public LogPosition apply(String destination) + { LogPosition logPosition = zooKeeperLogPositionManager.getLatestIndexBy(destination); - if (logPosition == null) { + if (logPosition == null) + { return nullPosition; - } else { + } else + { return logPosition; } } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/index/PeriodMixedLogPositionManager.java b/parse/src/main/java/com/alibaba/otter/canal/parse/index/PeriodMixedLogPositionManager.java index b1967163..4a12781b 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/index/PeriodMixedLogPositionManager.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/index/PeriodMixedLogPositionManager.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.parse.index; +import com.google.common.collect.MigrateMap; import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; @@ -43,13 +44,17 @@ public class PeriodMixedLogPositionManager extends MemoryLogPositionManager impl zooKeeperLogPositionManager.start(); } executor = Executors.newScheduledThreadPool(1); - positions = new MapMaker().makeComputingMap(new Function() { + positions = MigrateMap.makeComputingMap(new Function() + { - public LogPosition apply(String destination) { + public LogPosition apply(String destination) + { LogPosition logPosition = zooKeeperLogPositionManager.getLatestIndexBy(destination); - if (logPosition == null) { + if (logPosition == null) + { return nullPosition; - } else { + } else + { return logPosition; } } diff --git a/pom.xml b/pom.xml index 95dc3d86..d1d45299 100644 --- a/pom.xml +++ b/pom.xml @@ -1,311 +1,312 @@ - - 4.0.0 - com.alibaba.otter - canal - pom - canal module for otter ${project.version} - 1.0.20-SNAPSHOT - https://github.com/alibaba/canal - - org.sonatype.oss - oss-parent - 7 - - - - agapple - http://agapple.iteye.com - jianghang115@gmail.com - 8 - - - zavakid - http://www.zavakid.com - zava.kid@gmail.com - 8 - - - in355hz - http://in355hz.iteye.com - in355hz@gmail.com - 8 - - + + 4.0.0 + com.alibaba.otter + canal + pom + canal module for otter ${project.version} + 1.0.20-SNAPSHOT + https://github.com/alibaba/canal + + org.sonatype.oss + oss-parent + 7 + + + + agapple + http://agapple.iteye.com + jianghang115@gmail.com + 8 + + + zavakid + http://www.zavakid.com + zava.kid@gmail.com + 8 + + + in355hz + http://in355hz.iteye.com + in355hz@gmail.com + 8 + + - - - Apache License, Version 2.0 - http://www.apache.org/licenses/LICENSE-2.0 - - + + + Apache License, Version 2.0 + http://www.apache.org/licenses/LICENSE-2.0 + + - - git@github.com:alibaba/canal.git - scm:git:git@github.com:alibaba/canal.git - scm:git:git@github.com:alibaba/canal.git - - - - - central - http://repo1.maven.org/maven2 - - true - - - false - - - - java.net - http://download.java.net/maven/2/ - - true - - - false - - - - alibaba - http://code.alibabatech.com/mvn/releases/ - - true - - - false - - - - sonatype - sonatype - https://oss.sonatype.org/content/repositories/snapshots - - false - - - true - - - - sonatype-release - sonatype-release - https://oss.sonatype.org/service/local/repositories/releases/content - - false - - - true - - - - - - UTF-8 - - true - true - - 1.6 - 1.6 - UTF-8 - - - - common - meta - dbsync - filter - driver - parse - sink - store - protocol - instance - server - client - deployer - example - - - - - - org.springframework - spring - 2.5.6 - + + git@github.com:alibaba/canal.git + scm:git:git@github.com:alibaba/canal.git + scm:git:git@github.com:alibaba/canal.git + + + + + central + http://repo1.maven.org/maven2 + + true + + + false + + + + java.net + http://download.java.net/maven/2/ + + true + + + false + + + + alibaba + http://code.alibabatech.com/mvn/releases/ + + true + + + false + + + + sonatype + sonatype + https://oss.sonatype.org/content/repositories/snapshots + + false + + + true + + + + sonatype-release + sonatype-release + https://oss.sonatype.org/service/local/repositories/releases/content + + false + + + true + + + + + + UTF-8 + + true + true + + 1.6 + 1.6 + UTF-8 + + + + common + meta + dbsync + filter + driver + parse + sink + store + protocol + instance + server + client + deployer + example + + + + + + org.springframework + spring + 2.5.6 + - commons-lang - commons-lang - 2.6 - - - commons-io - commons-io - 2.4 - - - org.apache.zookeeper - zookeeper - 3.4.5 - - - log4j - log4j - - - org.slf4j - slf4j-log4j12 - - - org.slf4j - slf4j-api - - - jline - jline - - - - - com.github.sgroschupf - zkclient - 0.1 - - - com.alibaba - fastjson - 1.1.26 - - - com.google.guava - guava - r09 - - - com.googlecode.aviator - aviator - 2.2.1 - - - oro - oro - 2.0.8 - - - org.jboss.netty - netty - 3.2.5.Final - - - com.google.protobuf - protobuf-java - 2.4.1 - + commons-lang + commons-lang + 2.6 + + + commons-io + commons-io + 2.4 + + + org.apache.zookeeper + zookeeper + 3.4.5 + + + log4j + log4j + + + org.slf4j + slf4j-log4j12 + + + org.slf4j + slf4j-api + + + jline + jline + + + + + com.github.sgroschupf + zkclient + 0.1 + + + com.alibaba + fastjson + 1.1.26 + + + com.google.guava + guava + 18.0 + + + com.googlecode.aviator + aviator + 2.2.1 + + + oro + oro + 2.0.8 + + + org.jboss.netty + netty + 3.2.5.Final + + + com.google.protobuf + protobuf-java + 2.4.1 + - ch.qos.logback - logback-core - 1.0.6 - - - ch.qos.logback - logback-classic - 1.0.6 - - - org.slf4j - jcl-over-slf4j - 1.6.0 - - - org.slf4j - slf4j-api - 1.6.0 - - + ch.qos.logback + logback-core + 1.1.3 + - junit - junit - 4.5 - test - - - mysql - mysql-connector-java - 5.1.12 - test - - - - - - - - org.jvnet.wagon-svn - wagon-svn - 1.9 - - - org.apache.maven.wagon - wagon-http-shared - 1.0-beta-7 - - - - - org.apache.maven.plugins - maven-source-plugin + ch.qos.logback + logback-classic + 1.1.3 + + + org.slf4j + jcl-over-slf4j + 1.7.12 + + + org.slf4j + slf4j-api + 1.7.12 + + + + junit + junit + 4.12 + test + + + mysql + mysql-connector-java + 5.1.12 + test + + + + + + + + org.jvnet.wagon-svn + wagon-svn + 1.9 + + + org.apache.maven.wagon + wagon-http-shared + 1.0-beta-7 + + + + + org.apache.maven.plugins + maven-source-plugin 2.4 - - - attach-sources - - jar - - - - - - org.apache.maven.plugins - maven-compiler-plugin + + + attach-sources + + jar + + + + + + org.apache.maven.plugins + maven-compiler-plugin 3.2 - - ${java_source_version} - ${java_target_version} - ${file_encoding} - - - - org.apache.maven.plugins - maven-eclipse-plugin - 2.5.1 - - - - .settings/org.eclipse.core.resources.prefs - - =${file_encoding}${line.separator}]]> - - - - - - - org.apache.maven.plugins - maven-surefire-plugin - 2.5 - - - **/*Test.java - - - **/*NoRunTest.java - - - + + ${java_source_version} + ${java_target_version} + ${file_encoding} + + + + org.apache.maven.plugins + maven-eclipse-plugin + 2.5.1 + + + + .settings/org.eclipse.core.resources.prefs + + =${file_encoding}${line.separator}]]> + + + + + + + org.apache.maven.plugins + maven-surefire-plugin + 2.5 + + + **/*Test.java + + + **/*NoRunTest.java + + + - - src/main/java - src/test/java - - - src/main/resources - - **/* - - - **/.svn/ - - - - - - src/test/resources - - **/* - - - **/.svn/ - - - - - - - - sonatype-nexus-snapshots - Sonatype Nexus Snapshots - https://oss.sonatype.org/content/repositories/snapshots/ - - - sonatype-nexus-staging - Nexus Release Repository - https://oss.sonatype.org/service/local/staging/deploy/maven2/ - - + + src/main/java + src/test/java + + + src/main/resources + + **/* + + + **/.svn/ + + + + + + src/test/resources + + **/* + + + **/.svn/ + + + + + + + + sonatype-nexus-snapshots + Sonatype Nexus Snapshots + https://oss.sonatype.org/content/repositories/snapshots/ + + + sonatype-nexus-staging + Nexus Release Repository + https://oss.sonatype.org/service/local/staging/deploy/maven2/ + + diff --git a/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java b/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java index 604a54b6..c23ff139 100644 --- a/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java +++ b/server/src/main/java/com/alibaba/otter/canal/server/embedded/CanalServerWithEmbedded.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.server.embedded; +import com.google.common.collect.MigrateMap; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -47,9 +48,11 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C public void start() { super.start(); - canalInstances = new MapMaker().makeComputingMap(new Function() { + canalInstances = MigrateMap.makeComputingMap(new Function() + { - public CanalInstance apply(String destination) { + public CanalInstance apply(String destination) + { return canalInstanceGenerator.generate(destination); } });