adapter 远程配置

修改rdb热加载配置的问题
This commit is contained in:
machey
2019-01-05 01:42:27 +08:00
parent 3dec7c50ce
commit 779ebad890
11 changed files with 143 additions and 49 deletions
@@ -84,6 +84,10 @@ public class ESConfigMonitor {
// 加载配置文件
String configContent = MappingConfigsLoader
.loadConfig(adapterName + File.separator + file.getName());
if (configContent == null) {
onFileDelete(file);
return;
}
ESSyncConfig config = new Yaml().loadAs(configContent, ESSyncConfig.class);
config.validate();
if (esAdapter.getEsSyncConfig().containsKey(file.getName())) {
@@ -60,6 +60,10 @@ public class HbaseConfigMonitor {
try {
// 加载新增的配置文件
String configContent = MappingConfigsLoader.loadConfig(adapterName + File.separator + file.getName());
if (configContent == null) {
onFileDelete(file);
return;
}
MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class);
config.validate();
addConfigToCache(file, config);
@@ -1,5 +1,6 @@
package com.alibaba.otter.canal.adapter.launcher;
import com.alibaba.otter.canal.adapter.launcher.monitor.AdapterRemoteConfigMonitor;
import org.springframework.boot.Banner;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@@ -14,6 +15,21 @@ import org.springframework.boot.autoconfigure.SpringBootApplication;
public class CanalAdapterApplication {
public static void main(String[] args) {
// 加载远程配置
// String jdbcUrl = env.getProperty("canal.manager.jdbc.url");
// if (StringUtils.isNotEmpty(jdbcUrl)) {
// String jdbcUsername = env.getProperty("canal.manager.jdbc.username");
// String jdbcPassword = env.getProperty("canal.manager.jdbc.password");
// configMonitor = new AdapterRemoteConfigMonitor(jdbcUrl, jdbcUsername, jdbcPassword);
// configMonitor.loadRemoteConfig();
// configMonitor.loadRemoteAdapterConfigs();
// contextRefresher.refresh();
// configMonitor.start();
// }
// AdapterRemoteConfigMonitor aa = new AdapterRemoteConfigMonitor();
SpringApplication application = new SpringApplication(CanalAdapterApplication.class);
application.setBannerMode(Banner.Mode.OFF);
application.run(args);
@@ -0,0 +1,38 @@
package com.alibaba.otter.canal.adapter.launcher.config;
import javax.annotation.PostConstruct;
import com.alibaba.otter.canal.adapter.launcher.monitor.AdapterRemoteConfigMonitor;
import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.core.env.Environment;
public class BootstrapConfiguration {
private static final Logger logger = LoggerFactory.getLogger(BootstrapConfiguration.class);
@Autowired
private Environment env;
@PostConstruct
public void init() {
try {
// 加载远程配置
String jdbcUrl = env.getProperty("canal.manager.jdbc.url");
if (StringUtils.isNotEmpty(jdbcUrl)) {
String jdbcUsername = env.getProperty("canal.manager.jdbc.username");
String jdbcPassword = env.getProperty("canal.manager.jdbc.password");
AdapterRemoteConfigMonitor configMonitor = new AdapterRemoteConfigMonitor(jdbcUrl,
jdbcUsername,
jdbcPassword);
configMonitor.loadRemoteConfig();
configMonitor.loadRemoteAdapterConfigs();
configMonitor.start(); // 启动监听
}
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
}
}
@@ -57,18 +57,6 @@ public class CanalAdapterService {
return;
}
try {
// 加载远程配置
String jdbcUrl = env.getProperty("canal.manager.jdbc.url");
if (StringUtils.isNotEmpty(jdbcUrl)) {
String jdbcUsername = env.getProperty("canal.manager.jdbc.username");
String jdbcPassword = env.getProperty("canal.manager.jdbc.password");
configMonitor = new AdapterRemoteConfigMonitor(jdbcUrl, jdbcUsername, jdbcPassword);
configMonitor.loadRemoteConfig();
configMonitor.loadRemoteAdapterConfigs();
contextRefresher.refresh();
configMonitor.start();
}
logger.info("## start the canal client adapters.");
adapterLoader = new CanalAdapterLoader(adapterCanalConfig);
adapterLoader.init();
@@ -1,7 +1,9 @@
package com.alibaba.otter.canal.adapter.launcher.monitor;
import java.io.File;
import java.io.FileInputStream;
import java.io.FileWriter;
import java.net.URL;
import java.sql.*;
import java.util.ArrayList;
import java.util.HashMap;
@@ -11,12 +13,15 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import com.google.common.base.Joiner;
import com.google.common.collect.MapMaker;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.yaml.snakeyaml.Yaml;
import com.alibaba.otter.canal.client.adapter.support.Constant;
import com.alibaba.otter.canal.client.adapter.support.MappingConfigsLoader;
import com.alibaba.otter.canal.common.utils.NamedThreadFactory;
import com.google.common.base.Joiner;
import com.google.common.collect.MapMaker;
public class AdapterRemoteConfigMonitor {
@@ -39,6 +44,26 @@ public class AdapterRemoteConfigMonitor {
this.jdbcPassword = jdbcPassword;
}
public AdapterRemoteConfigMonitor(){
try {
File configFile = new File(".." + File.separator + Constant.CONF_DIR + File.separator + "bootstrap.yml");
if (!configFile.exists()) {
URL url = MappingConfigsLoader.class.getClassLoader().getResource("");
if (url != null) {
configFile = new File(url.getPath() + "bootstrap.yml");
}
}
if (configFile.exists()) {
try (FileInputStream fis = new FileInputStream(configFile)) {
Map a = new Yaml().load(fis);
a = a;
}
}
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
}
private Connection getConn() throws Exception {
if (conn == null || conn.isClosed()) {
Class.forName("com.mysql.jdbc.Driver");
@@ -124,7 +149,7 @@ public class AdapterRemoteConfigMonitor {
private Map<String, ConfigItem>[] getModifiedAdapterConfigs() {
Map<String, ConfigItem>[] res = new Map[2];
Map<String, ConfigItem> remoteConfigStatus = new HashMap<>();
String sql = "select id, category, name, modified_time from canal_instance_config";
String sql = "select id, category, name, modified_time from canal_adapter_config";
try (Statement stmt = getConn().createStatement(); ResultSet rs = stmt.executeQuery(sql)) {
while (rs.next()) {
ConfigItem configItem = new ConfigItem();
@@ -132,7 +157,7 @@ public class AdapterRemoteConfigMonitor {
configItem.setCategory(rs.getString("category"));
configItem.setName(rs.getString("name"));
configItem.setModifiedTime(rs.getTimestamp("modified_time").getTime());
remoteConfigStatus.put(configItem.getName(), configItem);
remoteConfigStatus.put(configItem.getCategory() + "/" + configItem.getName(), configItem);
}
} catch (Exception e) {
logger.error(e.getMessage(), e);
@@ -157,7 +182,7 @@ public class AdapterRemoteConfigMonitor {
}
if (!changedIds.isEmpty()) {
Map<String, ConfigItem> changedInstanceConfig = new HashMap<>();
String contentsSql = "select id, name, content, modified_time from canal_instance_config where id in ("
String contentsSql = "select id, category, name, content, modified_time from canal_adapter_config where id in ("
+ Joiner.on(",").join(changedIds) + ")";
try (Statement stmt = getConn().createStatement(); ResultSet rs = stmt.executeQuery(contentsSql)) {
while (rs.next()) {
@@ -182,11 +207,11 @@ public class AdapterRemoteConfigMonitor {
}
Map<String, ConfigItem> removedInstanceConfig = new HashMap<>();
for (String name : remoteAdapterConfigs.keySet()) {
if (!remoteConfigStatus.containsKey(name)) {
for (ConfigItem configItem : remoteAdapterConfigs.values()) {
if (!remoteConfigStatus.containsKey(configItem.getCategory() + "/" + configItem.getName())) {
// 删除
remoteAdapterConfigs.remove(name);
removedInstanceConfig.put(name, null);
remoteAdapterConfigs.remove(configItem.getCategory() + "/" + configItem.getName());
removedInstanceConfig.put(configItem.getCategory() + "/" + configItem.getName(), null);
}
}
res[1] = removedInstanceConfig.isEmpty() ? null : removedInstanceConfig;
@@ -0,0 +1,3 @@
# Bootstrap Configuration
org.springframework.cloud.bootstrap.BootstrapConfiguration=\
com.alibaba.otter.canal.adapter.launcher.config.BootstrapConfiguration
@@ -12,11 +12,6 @@ spring:
time-zone: GMT+8
default-property-inclusion: non_null
#canal.manager.jdbc:
# url: jdbc:mysql://127.0.0.1:3306/canal_manager?useUnicode=true&characterEncoding=UTF-8
# username: root
# password: 121212
canal.conf:
canalServerHost: 127.0.0.1:11111
# zookeeperHosts: slave1:2181
@@ -44,7 +39,7 @@ canal.conf:
# key: mysql1
# properties:
# jdbc.driverClassName: com.mysql.jdbc.Driver
# jdbc.url: jdbc:mysql://192.168.0.36/mytest?useUnicode=true
# jdbc.url: jdbc:mysql://127.0.0.1:3306/mytest2?useUnicode=true
# jdbc.username: root
# jdbc.password: 121212
# - name: rdb
@@ -0,0 +1,6 @@
canal:
manager:
jdbc:
url: jdbc:mysql://127.0.0.1:3306/canal_manager?useUnicode=true&characterEncoding=UTF-8
username: root
password: 121212
@@ -4,7 +4,6 @@ import java.io.File;
import java.util.HashMap;
import java.util.Map;
import com.alibaba.otter.canal.client.adapter.rdb.config.MirrorDbConfig;
import org.apache.commons.io.filefilter.FileFilterUtils;
import org.apache.commons.io.monitor.FileAlterationListenerAdaptor;
import org.apache.commons.io.monitor.FileAlterationMonitor;
@@ -16,6 +15,7 @@ import org.yaml.snakeyaml.Yaml;
import com.alibaba.otter.canal.client.adapter.rdb.RdbAdapter;
import com.alibaba.otter.canal.client.adapter.rdb.config.MappingConfig;
import com.alibaba.otter.canal.client.adapter.rdb.config.MirrorDbConfig;
import com.alibaba.otter.canal.client.adapter.support.MappingConfigsLoader;
import com.alibaba.otter.canal.client.adapter.support.Util;
@@ -86,6 +86,10 @@ public class RdbConfigMonitor {
// 加载配置文件
String configContent = MappingConfigsLoader
.loadConfig(adapterName + File.separator + file.getName());
if (configContent == null) {
onFileDelete(file);
return;
}
MappingConfig config = new Yaml().loadAs(configContent, MappingConfig.class);
config.validate();
if ((key == null && config.getOuterAdapterKey() == null)
@@ -121,34 +125,44 @@ public class RdbConfigMonitor {
}
private void addConfigToCache(File file, MappingConfig mappingConfig) {
if (mappingConfig == null || mappingConfig.getDbMapping() == null) {
return;
}
rdbAdapter.getRdbMapping().put(file.getName(), mappingConfig);
Map<String, MappingConfig> configMap = rdbAdapter.getMappingConfigCache()
.computeIfAbsent(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
+ mappingConfig.getDbMapping().getDatabase() + "."
+ mappingConfig.getDbMapping().getTable(),
k1 -> new HashMap<>());
configMap.put(file.getName(), mappingConfig);
Map<String, MirrorDbConfig> mirrorDbConfigCache = rdbAdapter.getMirrorDbConfigCache();
mirrorDbConfigCache.put(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
+ mappingConfig.getDbMapping().getDatabase(),
MirrorDbConfig.create(file.getName(), mappingConfig));
if (!mappingConfig.getDbMapping().getMirrorDb()) {
Map<String, MappingConfig> configMap = rdbAdapter.getMappingConfigCache()
.computeIfAbsent(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
+ mappingConfig.getDbMapping().getDatabase() + "."
+ mappingConfig.getDbMapping().getTable(),
k1 -> new HashMap<>());
configMap.put(file.getName(), mappingConfig);
} else {
Map<String, MirrorDbConfig> mirrorDbConfigCache = rdbAdapter.getMirrorDbConfigCache();
mirrorDbConfigCache.put(StringUtils.trimToEmpty(mappingConfig.getDestination()) + "."
+ mappingConfig.getDbMapping().getDatabase(),
MirrorDbConfig.create(file.getName(), mappingConfig));
}
}
private void deleteConfigFromCache(File file) {
MappingConfig mappingConfig = rdbAdapter.getRdbMapping().remove(file.getName());
rdbAdapter.getRdbMapping().remove(file.getName());
for (Map<String, MappingConfig> configMap : rdbAdapter.getMappingConfigCache().values()) {
if (configMap != null) {
configMap.remove(file.getName());
}
if (mappingConfig == null || mappingConfig.getDbMapping() == null) {
return;
}
rdbAdapter.getMirrorDbConfigCache().forEach((key, mirrorDbConfig) -> {
if (mirrorDbConfig.getFileName().equals(file.getName())) {
rdbAdapter.getMirrorDbConfigCache().remove(key);
if (!mappingConfig.getDbMapping().getMirrorDb()) {
for (Map<String, MappingConfig> configMap : rdbAdapter.getMappingConfigCache().values()) {
if (configMap != null) {
configMap.remove(file.getName());
}
}
});
} else {
rdbAdapter.getMirrorDbConfigCache().forEach((key, mirrorDbConfig) -> {
if (mirrorDbConfig.getFileName().equals(file.getName())) {
rdbAdapter.getMirrorDbConfigCache().remove(key);
}
});
}
}
}
+1
View File
@@ -28,6 +28,7 @@ CREATE TABLE `canal_adapter_config` (
) ENGINE=InnoDB AUTO_INCREMENT=2 DEFAULT CHARSET=utf8mb4 ;
INSERT INTO `canal_config` VALUES ('1', 'canal.properties', '#################################################\r\n######### common argument ############# \r\n#################################################\r\n#canal.manager.jdbc.url=jdbc:mysql://127.0.0.1:3306/canal_manager?useUnicode=true&characterEncoding=UTF-8\r\n#canal.manager.jdbc.username=root\r\n#canal.manager.jdbc.password=121212\r\ncanal.id = 1\r\ncanal.ip =\r\ncanal.port = 11111\r\ncanal.metrics.pull.port = 11112\r\ncanal.zkServers =\r\n# flush data to zk\r\ncanal.zookeeper.flush.period = 1000\r\ncanal.withoutNetty = false\r\n# tcp, kafka, RocketMQ\r\ncanal.serverMode = tcp\r\n# flush meta cursor/parse position to file\r\ncanal.file.data.dir = ${canal.conf.dir}\r\ncanal.file.flush.period = 1000\r\n## memory store RingBuffer size, should be Math.pow(2,n)\r\ncanal.instance.memory.buffer.size = 16384\r\n## memory store RingBuffer used memory unit size , default 1kb\r\ncanal.instance.memory.buffer.memunit = 1024 \r\n## meory store gets mode used MEMSIZE or ITEMSIZE\r\ncanal.instance.memory.batch.mode = MEMSIZE\r\ncanal.instance.memory.rawEntry = true\r\n\r\n## detecing config\r\ncanal.instance.detecting.enable = false\r\n#canal.instance.detecting.sql = insert into retl.xdual values(1,now()) on duplicate key update x=now()\r\ncanal.instance.detecting.sql = select 1\r\ncanal.instance.detecting.interval.time = 3\r\ncanal.instance.detecting.retry.threshold = 3\r\ncanal.instance.detecting.heartbeatHaEnable = false\r\n\r\n# support maximum transaction size, more than the size of the transaction will be cut into multiple transactions delivery\r\ncanal.instance.transaction.size = 1024\r\n# mysql fallback connected to new master should fallback times\r\ncanal.instance.fallbackIntervalInSeconds = 60\r\n\r\n# network config\r\ncanal.instance.network.receiveBufferSize = 16384\r\ncanal.instance.network.sendBufferSize = 16384\r\ncanal.instance.network.soTimeout = 30\r\n\r\n# binlog filter config\r\ncanal.instance.filter.druid.ddl = true\r\ncanal.instance.filter.query.dcl = false\r\ncanal.instance.filter.query.dml = false\r\ncanal.instance.filter.query.ddl = false\r\ncanal.instance.filter.table.error = false\r\ncanal.instance.filter.rows = false\r\ncanal.instance.filter.transaction.entry = false\r\n\r\n# binlog format/image check\r\ncanal.instance.binlog.format = ROW,STATEMENT,MIXED \r\ncanal.instance.binlog.image = FULL,MINIMAL,NOBLOB\r\n\r\n# binlog ddl isolation\r\ncanal.instance.get.ddl.isolation = false\r\n\r\n# parallel parser config\r\ncanal.instance.parser.parallel = true\r\n## concurrent thread number, default 60% available processors, suggest not to exceed Runtime.getRuntime().availableProcessors()\r\n#canal.instance.parser.parallelThreadSize = 16\r\n## disruptor ringbuffer size, must be power of 2\r\ncanal.instance.parser.parallelBufferSize = 256\r\n\r\n# table meta tsdb info\r\ncanal.instance.tsdb.enable = true\r\ncanal.instance.tsdb.dir = ${canal.file.data.dir:../conf}/${canal.instance.destination:}\r\ncanal.instance.tsdb.url = jdbc:h2:${canal.instance.tsdb.dir}/h2;CACHE_SIZE=1000;MODE=MYSQL;\r\ncanal.instance.tsdb.dbUsername = canal\r\ncanal.instance.tsdb.dbPassword = canal\r\n# dump snapshot interval, default 24 hour\r\ncanal.instance.tsdb.snapshot.interval = 24\r\n# purge snapshot expire , default 360 hour(15 days)\r\ncanal.instance.tsdb.snapshot.expire = 360\r\n\r\n# aliyun ak/sk , support rds/mq\r\ncanal.aliyun.accesskey =\r\ncanal.aliyun.secretkey =\r\n\r\n#################################################\r\n######### destinations ############# \r\n#################################################\r\ncanal.destinations = example\r\n# conf root dir\r\ncanal.conf.dir = ../conf\r\n# auto scan instance dir add/remove and start/stop instance\r\ncanal.auto.scan = true\r\ncanal.auto.scan.interval = 5\r\n\r\ncanal.instance.tsdb.spring.xml = classpath:spring/tsdb/h2-tsdb.xml\r\n#canal.instance.tsdb.spring.xml = classpath:spring/tsdb/mysql-tsdb.xml\r\n\r\ncanal.instance.global.mode = spring\r\ncanal.instance.global.lazy = false\r\n#canal.instance.global.manager.address = 127.0.0.1:1099\r\n#canal.instance.global.spring.xml = classpath:spring/memory-instance.xml\r\ncanal.instance.global.spring.xml = classpath:spring/file-instance.xml\r\n#canal.instance.global.spring.xml = classpath:spring/default-instance.xml\r\n\r\n##################################################\r\n######### MQ #############\r\n##################################################\r\ncanal.mq.servers = 127.0.0.1:6667\r\ncanal.mq.retries = 0\r\ncanal.mq.batchSize = 16384\r\ncanal.mq.maxRequestSize = 1048576\r\ncanal.mq.lingerMs = 1\r\ncanal.mq.bufferMemory = 33554432\r\ncanal.mq.canalBatchSize = 50\r\ncanal.mq.canalGetTimeout = 100\r\ncanal.mq.flatMessage = true\r\ncanal.mq.compressionType = none\r\ncanal.mq.acks = all\r\n', '2018-12-30 16:46:00');
INSERT INTO `canal_config` VALUES ('2', 'application.yml', 'server:\n port: 8081\nlogging:\n level:\n org.springframework: WARN\n com.alibaba.otter.canal.client.adapter.hbase: DEBUG\n com.alibaba.otter.canal.client.adapter.es: DEBUG\n com.alibaba.otter.canal.client.adapter.rdb: DEBUG\nspring:\n jackson:\n date-format: yyyy-MM-dd HH:mm:ss\n time-zone: GMT+8\n default-property-inclusion: non_null\n\ncanal.conf:\n canalServerHost: 127.0.0.1:11111\n# zookeeperHosts: slave1:2181\n# mqServers: 127.0.0.1:9092 #or rocketmq\n# flatMessage: true\n batchSize: 500\n syncBatchSize: 1000\n retries: 0\n timeout:\n accessKey:\n secretKey:\n mode: tcp # kafka rocketMQ\n# srcDataSources:\n# defaultDS:\n# url: jdbc:mysql://127.0.0.1:3306/mytest?useUnicode=true\n# username: root\n# password: 121212\n canalAdapters:\n - instance: example # canal instance Name or mq topic name\n groups:\n - groupId: g1\n outerAdapters:\n - name: logger\n# - name: rdb\n# key: mysql1\n# properties:\n# jdbc.driverClassName: com.mysql.jdbc.Driver\n# jdbc.url: jdbc:mysql://127.0.0.1:3306/mytest2?useUnicode=true\n# jdbc.username: root\n# jdbc.password: 121212\n# - name: rdb\n# key: oracle1\n# properties:\n# jdbc.driverClassName: oracle.jdbc.OracleDriver\n# jdbc.url: jdbc:oracle:thin:@localhost:49161:XE\n# jdbc.username: mytest\n# jdbc.password: m121212\n# - name: rdb\n# key: postgres1\n# properties:\n# jdbc.driverClassName: org.postgresql.Driver\n# jdbc.url: jdbc:postgresql://localhost:5432/postgres\n# jdbc.username: postgres\n# jdbc.password: 121212\n# threads: 1\n# commitSize: 3000\n# - name: hbase\n# properties:\n# hbase.zookeeper.quorum: 127.0.0.1\n# hbase.zookeeper.property.clientPort: 2181\n# zookeeper.znode.parent: /hbase\n# - name: es\n# hosts: 127.0.0.1:9300\n# properties:\n# cluster.name: elasticsearch\n', '2019-01-05 01:37:43');
INSERT INTO `canal_instance_config` VALUES ('1', 'example', '#################################################\r\n## mysql serverId , v1.0.26+ will autoGen\r\n# canal.instance.mysql.slaveId=0\r\n\r\n# enable gtid use true/false\r\ncanal.instance.gtidon=false\r\n\r\n# position info\r\ncanal.instance.master.address=127.0.0.1:3306\r\ncanal.instance.master.journal.name=\r\ncanal.instance.master.position=\r\ncanal.instance.master.timestamp=\r\ncanal.instance.master.gtid=\r\n\r\n# rds oss binlog\r\ncanal.instance.rds.accesskey=\r\ncanal.instance.rds.secretkey=\r\ncanal.instance.rds.instanceId=\r\n\r\n# table meta tsdb info\r\ncanal.instance.tsdb.enable=true\r\n#canal.instance.tsdb.url=jdbc:mysql://127.0.0.1:3306/canal_tsdb\r\n#canal.instance.tsdb.dbUsername=canal\r\n#canal.instance.tsdb.dbPassword=canal\r\n\r\n#canal.instance.standby.address =\r\n#canal.instance.standby.journal.name =\r\n#canal.instance.standby.position =\r\n#canal.instance.standby.timestamp =\r\n#canal.instance.standby.gtid=\r\n\r\n# username/password\r\ncanal.instance.dbUsername=canal\r\ncanal.instance.dbPassword=canal\r\ncanal.instance.connectionCharset = UTF-8\r\ncanal.instance.defaultDatabaseName =test\r\n# enable druid Decrypt database password\r\ncanal.instance.enableDruid=false\r\n#canal.instance.pwdPublicKey=MFwwDQYJKoZIhvcNAQEBBQADSwAwSAJBALK4BUxdDltRRE5/zXpVEVPUgunvscYFtEip3pmLlhrWpacX7y7GCMo2/JM6LeHmiiNdH1FWgGCpUfircSwlWKUCAwEAAQ==\r\n\r\n# table regex\r\ncanal.instance.filter.regex=.*\\\\..*\r\n# table black regex\r\ncanal.instance.filter.black.regex=\r\n\r\n# mq config\r\ncanal.mq.topic=example\r\n# 动态topic, 需mq支持动态创建topic\r\n#canal.mq.dynamicTopic=.*,mytest\\\\..*,mytest2.user\r\ncanal.mq.partition=0\r\n# hash partition config\r\n#canal.mq.partitionsNum=3\r\n#canal.mq.partitionHash=test.table:id^name,.*\\\\..*\r\n#################################################\r\n', '2018-12-30 16:46:16');