diff --git a/.travis.yml b/.travis.yml index d0332f54..005c4404 100644 --- a/.travis.yml +++ b/.travis.yml @@ -1,4 +1,5 @@ language: java +sudo: false # faster builds jdk: - oraclejdk12 diff --git a/admin/admin-ui/src/api/canalInstance.js b/admin/admin-ui/src/api/canalInstance.js index 2d7f2ac2..dbb81c7a 100644 --- a/admin/admin-ui/src/api/canalInstance.js +++ b/admin/admin-ui/src/api/canalInstance.js @@ -59,13 +59,6 @@ export function instanceLog(id, nodeId) { }) } -export function instanceMeta(id, nodeId) { - return request({ - url: '/canal/instance/meta/' + id + '/' + nodeId, - method: 'get' - }) -} - export function instanceStatus(id, option) { return request({ url: '/canal/instance/status/' + id + '?option=' + option, diff --git a/admin/admin-ui/src/router/index.js b/admin/admin-ui/src/router/index.js index d1ac3b56..269022e7 100644 --- a/admin/admin-ui/src/router/index.js +++ b/admin/admin-ui/src/router/index.js @@ -128,13 +128,6 @@ export const constantRoutes = [ component: () => import('@/views/canalServer/CanalInstanceLogDetail'), meta: { title: 'Instance 日志' }, hidden: true - }, - { - path: 'canalInstance/meta', - name: 'Instance meta', - component: () => import('@/views/canalServer/CanalInstanceMetaDetail'), - meta: { title: 'Instance Meta' }, - hidden: true } ] }, diff --git a/admin/admin-ui/src/views/canalServer/CanalInstance.vue b/admin/admin-ui/src/views/canalServer/CanalInstance.vue index 918762f4..0ec9e56b 100644 --- a/admin/admin-ui/src/views/canalServer/CanalInstance.vue +++ b/admin/admin-ui/src/views/canalServer/CanalInstance.vue @@ -64,7 +64,6 @@ 启动 停止 日志 - meta @@ -118,7 +117,6 @@ export default { } }, created() { - this.listQuery.name = this.$route.query.name getClustersAndServers().then((res) => { this.options = res.data }) @@ -224,13 +222,6 @@ export default { return } this.$router.push('canalInstance/log?id=' + row.id + '&nodeId=' + row.nodeServer.id) - }, - handleMeta(row) { - if (row.nodeId === null) { - this.$message({ message: '当前Instance不是启动状态,无法查看meta', type: 'warning' }) - return - } - this.$router.push('canalInstance/meta?id=' + row.id + '&nodeId=' + row.nodeServer.id) } } } diff --git a/admin/admin-ui/src/views/canalServer/CanalInstanceMetaDetail.vue b/admin/admin-ui/src/views/canalServer/CanalInstanceMetaDetail.vue deleted file mode 100644 index 19a09aef..00000000 --- a/admin/admin-ui/src/views/canalServer/CanalInstanceMetaDetail.vue +++ /dev/null @@ -1,53 +0,0 @@ - - - - - - {{ form.instance }}.meta - 刷新 - 返回 - - - - - - - - - - - diff --git a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/AdminConnector.java b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/AdminConnector.java index 1b1368a0..d06b29cb 100644 --- a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/AdminConnector.java +++ b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/AdminConnector.java @@ -128,12 +128,4 @@ public interface AdminConnector { */ String instanceLog(String destination, String fileName, int lines); - /** - * meta - * @param destination - * @param fileName - * @return - */ - String instanceMeta(String destination, String fileName); - } diff --git a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/SimpleAdminConnector.java b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/SimpleAdminConnector.java index eef3c3e2..d3677280 100644 --- a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/SimpleAdminConnector.java +++ b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/connector/SimpleAdminConnector.java @@ -206,11 +206,6 @@ public class SimpleAdminConnector implements AdminConnector { return doLogAdmin("instance", "file", destination, fileName, lines); } - @Override - public String instanceMeta(final String destination, final String fileName) { - return doLogAdmin("meta", "file", destination, fileName,100); - } - // ==================== helper method ==================== private String doServerAdmin(String action) { diff --git a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/controller/CanalInstanceController.java b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/controller/CanalInstanceController.java index 89e73802..e518557b 100644 --- a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/controller/CanalInstanceController.java +++ b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/controller/CanalInstanceController.java @@ -161,20 +161,6 @@ public class CanalInstanceController { return BaseModel.getInstance(canalInstanceConfigService.remoteInstanceLog(id, nodeId)); } - /** - * 获取实例meta信息 - * - * @param id - * @param nodeId - * @param env - * @return - */ - @GetMapping(value = "/instance/meta/{id}/{nodeId}") - public BaseModel> meta(@PathVariable Long id, @PathVariable Long nodeId, - @PathVariable String env) { - return BaseModel.getInstance(canalInstanceConfigService.remoteInstanceMeta(id, nodeId)); - } - /** * 通过Server id获取所有活动的Instance * diff --git a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/CanalInstanceService.java b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/CanalInstanceService.java index 0ed45378..fedd188c 100644 --- a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/CanalInstanceService.java +++ b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/CanalInstanceService.java @@ -26,8 +26,6 @@ public interface CanalInstanceService { Map remoteInstanceLog(Long id, Long nodeId); - Map remoteInstanceMeta(Long id, Long nodeId); - boolean remoteOperation(Long id, Long nodeId, String option); boolean instanceOperation(Long id, String option); diff --git a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/impl/CanalInstanceServiceImpl.java b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/impl/CanalInstanceServiceImpl.java index 5e0408d3..4d1a9d83 100644 --- a/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/impl/CanalInstanceServiceImpl.java +++ b/admin/admin-web/src/main/java/com/alibaba/otter/canal/admin/service/impl/CanalInstanceServiceImpl.java @@ -1,10 +1,8 @@ package com.alibaba.otter.canal.admin.service.impl; -import com.alibaba.otter.canal.common.zookeeper.ZkClientx; import io.ebean.Query; import java.security.NoSuchAlgorithmException; -import java.text.MessageFormat; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -258,33 +256,6 @@ public class CanalInstanceServiceImpl implements CanalInstanceService { return result; } - @Override - public Map remoteInstanceMeta(final Long id, final Long nodeId) { - Map result = new HashMap<>(); - - NodeServer nodeServer = NodeServer.find.byId(nodeId); - if (nodeServer == null) { - return result; - } - CanalInstanceConfig canalInstanceConfig = CanalInstanceConfig.find.byId(id); - if (canalInstanceConfig == null) { - return result; - } - String meta; - if (nodeServer.getCanalCluster() != null) { - ZkClientx zkClientx = ZkClientx.getZkClient(nodeServer.getCanalCluster().getZkHosts()); - String zkPath = MessageFormat.format("/{0}/{1}/{2}/{3}/{4}/{5}", "otter", "canal", "destinations", canalInstanceConfig.getName(), "1001", "cursor"); - meta = new String((byte[]) zkClientx.readData(zkPath)); - } else { - meta = SimpleAdminConnectors.execute(nodeServer.getIp(), - nodeServer.getAdminPort(), - adminConnector -> adminConnector.instanceMeta(canalInstanceConfig.getName(), "meta.dat")); - } - result.put("instance", canalInstanceConfig.getName()); - result.put("meta", meta); - return result; - } - public boolean remoteOperation(Long id, Long nodeId, String option) { NodeServer nodeServer = null; if ("start".equals(option)) { diff --git a/admin/admin-web/src/main/resources/canal-template.properties b/admin/admin-web/src/main/resources/canal-template.properties index 381c913f..8a104e68 100644 --- a/admin/admin-web/src/main/resources/canal-template.properties +++ b/admin/admin-web/src/main/resources/canal-template.properties @@ -8,11 +8,11 @@ canal.register.ip = canal.port = 11111 canal.metrics.pull.port = 11112 # canal instance user/passwd -canal.user = canal -canal.passwd = E3619321C1A937C46A0D8BD1DAC39F93B27D4458 +# canal.user = canal +# canal.passwd = E3619321C1A937C46A0D8BD1DAC39F93B27D4458 # canal admin config -canal.admin.manager = 127.0.0.1:8089 +#canal.admin.manager = 127.0.0.1:8089 canal.admin.port = 11110 canal.admin.user = admin canal.admin.passwd = 4ACFE3202A5FF5CF467898FC58AAB1D615029441 @@ -21,7 +21,7 @@ canal.zkServers = # flush data to zk canal.zookeeper.flush.period = 1000 canal.withoutNetty = false -# tcp, kafka, RocketMQ +# tcp, kafka, rocketMQ, rabbitMQ canal.serverMode = tcp # flush meta cursor/parse position to file canal.file.data.dir = ${canal.conf.dir} @@ -86,14 +86,10 @@ canal.instance.tsdb.snapshot.interval = 24 # purge snapshot expire , default 360 hour(15 days) canal.instance.tsdb.snapshot.expire = 360 -# aliyun ak/sk , support rds/mq -canal.aliyun.accessKey = -canal.aliyun.secretKey = - ################################################# ######### destinations ############# ################################################# -canal.destinations = +canal.destinations = # conf root dir canal.conf.dir = ../conf # auto scan instance dir add/remove and start/stop instance @@ -111,29 +107,56 @@ canal.instance.global.spring.xml = classpath:spring/file-instance.xml #canal.instance.global.spring.xml = classpath:spring/default-instance.xml ################################################## -######### MQ ############# +######### MQ Properties ############# ################################################## -canal.mq.servers = 127.0.0.1:6667 -canal.mq.retries = 0 -canal.mq.batchSize = 16384 -canal.mq.maxRequestSize = 1048576 -canal.mq.lingerMs = 100 -canal.mq.bufferMemory = 33554432 +# aliyun ak/sk , support rds/mq +canal.aliyun.accessKey = +canal.aliyun.secretKey = +canal.aliyun.uid= + +canal.mq.flatMessage = true canal.mq.canalBatchSize = 50 canal.mq.canalGetTimeout = 100 -canal.mq.flatMessage = true -canal.mq.compressionType = none -canal.mq.acks = all -#canal.mq.properties. = -canal.mq.producerGroup = test # Set this value to "cloud", if you want open message trace feature in aliyun. canal.mq.accessChannel = local -# aliyun mq namespace -#canal.mq.namespace = + +canal.mq.database.hash = true +canal.mq.send.thread.size = 30 +canal.mq.build.thread.size = 8 ################################################## -######### Kafka Kerberos Info ############# +######### Kafka ############# ################################################## -canal.mq.kafka.kerberos.enable = false -canal.mq.kafka.kerberos.krb5FilePath = "../conf/kerberos/krb5.conf" -canal.mq.kafka.kerberos.jaasFilePath = "../conf/kerberos/jaas.conf" +kafka.bootstrap.servers = 127.0.0.1:6667 +kafka.acks = all +kafka.compression.type = none +kafka.batch.size = 16384 +kafka.linger.ms = 1 +kafka.max.request.size = 1048576 +kafka.buffer.memory = 33554432 +kafka.max.in.flight.requests.per.connection = 1 +kafka.retries = 0 + +kafka.kerberos.enable = false +kafka.kerberos.krb5.file = "../conf/kerberos/krb5.conf" +kafka.kerberos.jaas.file = "../conf/kerberos/jaas.conf" + +################################################## +######### RocketMQ ############# +################################################## +rocketmq.producer.group = test +rocketmq.enable.message.trace = false +rocketmq.customized.trace.topic = +rocketmq.namespace = +rocketmq.namesrv.addr = 127.0.0.1:9876 +rocketmq.retry.times.when.send.failed = 0 +rocketmq.vip.channel.enabled = false + +################################################## +######### RabbitMQ ############# +################################################## +rabbitmq.host = +rabbitmq.virtual.host = +rabbitmq.exchange = +rabbitmq.username = +rabbitmq.password = \ No newline at end of file diff --git a/admin/pom.xml b/admin/pom.xml index 5f8b45df..f92fe477 100644 --- a/admin/pom.xml +++ b/admin/pom.xml @@ -69,7 +69,7 @@ mysql mysql-connector-java - 8.0.17 + 5.1.48 com.github.ben-manes.caffeine diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/config/bind/PropertiesConfigurationFactory.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/config/bind/PropertiesConfigurationFactory.java index f44c9a5a..76a97c0b 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/config/bind/PropertiesConfigurationFactory.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/config/bind/PropertiesConfigurationFactory.java @@ -282,7 +282,7 @@ public class PropertiesConfigurationFactory implements FactoryBean, Applic PropertyDescriptor[] descriptors = BeanUtils.getPropertyDescriptors(this.target.getClass()); for (PropertyDescriptor descriptor : descriptors) { String name = descriptor.getName(); - if (!"class".equals(name)) { + if (!name.equals("class")) { RelaxedNames relaxedNames = RelaxedNames.forCamelCase(name); if (prefixes == null) { for (String relaxedName : relaxedNames) { diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java index a9a27960..4dd061c6 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java @@ -78,7 +78,7 @@ public class JdbcTypeUtil { public static Object typeConvert(String tableName ,String columnName, String value, int sqlType, String mysqlType) { if (value == null - || ("".equals(value) && !(isText(mysqlType) || sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.LONGVARCHAR))) { + || (value.equals("") && !(isText(mysqlType) || sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.LONGVARCHAR))) { return null; } diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java index 8d12debc..d0b42951 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java @@ -38,6 +38,7 @@ public class MessageUtil { CanalEntry.EventType eventType = rowChange.getEventType(); final Dml dml = new Dml(); + dml.setIsDdl(rowChange.getIsDdl()); dml.setDestination(destination); dml.setGroupId(groupId); dml.setDatabase(entry.getHeader().getSchemaName()); diff --git a/client-adapter/escore/src/main/java/com/alibaba/otter/canal/client/adapter/es/core/service/ESSyncService.java b/client-adapter/escore/src/main/java/com/alibaba/otter/canal/client/adapter/es/core/service/ESSyncService.java index 7b4b86a3..4ca23af2 100644 --- a/client-adapter/escore/src/main/java/com/alibaba/otter/canal/client/adapter/es/core/service/ESSyncService.java +++ b/client-adapter/escore/src/main/java/com/alibaba/otter/canal/client/adapter/es/core/service/ESSyncService.java @@ -95,11 +95,11 @@ public class ESSyncService { long begin = System.currentTimeMillis(); String type = dml.getType(); - if (type != null && "INSERT".equalsIgnoreCase(type)) { + if (type != null && type.equalsIgnoreCase("INSERT")) { insert(config, dml); - } else if (type != null && "UPDATE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("UPDATE")) { update(config, dml); - } else if (type != null && "DELETE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("DELETE")) { delete(config, dml); } else { return; diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java index 4940a58c..21ac327c 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java @@ -312,7 +312,7 @@ public class MappingConfig implements AdapterConfig { columnItem.setRowKey(true); rowKeyColumn = columnItem; } else { - if (field == null || "".equals(field)) { + if (field == null || field.equals("")) { columnItem.setFamily(family); columnItem.setQualifier(columnField.getKey()); } else { diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java index b4464cab..e0dc1279 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java @@ -31,11 +31,11 @@ public class HbaseSyncService { public void sync(MappingConfig config, Dml dml) { if (config != null) { String type = dml.getType(); - if (type != null && "INSERT".equalsIgnoreCase(type)) { + if (type != null && type.equalsIgnoreCase("INSERT")) { insert(config, dml); - } else if (type != null && "UPDATE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("UPDATE")) { update(config, dml); - } else if (type != null && "DELETE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("DELETE")) { delete(config, dml); } if (logger.isDebugEnabled()) { diff --git a/client-adapter/kudu/pom.xml b/client-adapter/kudu/pom.xml index 10c87858..3024a638 100644 --- a/client-adapter/kudu/pom.xml +++ b/client-adapter/kudu/pom.xml @@ -34,7 +34,7 @@ mysql mysql-connector-java - 5.1.47 + 5.1.48 test @@ -84,4 +84,4 @@ - \ No newline at end of file + diff --git a/client-adapter/kudu/src/main/java/com/alibaba/otter/canal/client/adapter/kudu/service/KuduSyncService.java b/client-adapter/kudu/src/main/java/com/alibaba/otter/canal/client/adapter/kudu/service/KuduSyncService.java index 4f5ca799..d388b410 100644 --- a/client-adapter/kudu/src/main/java/com/alibaba/otter/canal/client/adapter/kudu/service/KuduSyncService.java +++ b/client-adapter/kudu/src/main/java/com/alibaba/otter/canal/client/adapter/kudu/service/KuduSyncService.java @@ -47,11 +47,11 @@ public class KuduSyncService { public void sync(KuduMappingConfig config, Dml dml) { if (config != null) { String type = dml.getType(); - if (type != null && "INSERT".equalsIgnoreCase(type)) { + if (type != null && type.equalsIgnoreCase("INSERT")) { insert(config, dml); - } else if (type != null && "UPDATE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("UPDATE")) { upsert(config, dml); - } else if (type != null && "DELETE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("DELETE")) { delete(config, dml); } if (logger.isDebugEnabled()) { diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java index 0ec60392..ecef9199 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/config/CuratorClient.java @@ -3,6 +3,7 @@ package com.alibaba.otter.canal.adapter.launcher.config; import javax.annotation.PostConstruct; import javax.annotation.Resource; +import org.apache.commons.lang.StringUtils; import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; import org.apache.curator.retry.ExponentialBackoffRetry; @@ -24,7 +25,7 @@ public class CuratorClient { @PostConstruct public void init() { - if (adapterCanalConfig.getZookeeperHosts() != null) { + if (StringUtils.isNotEmpty(adapterCanalConfig.getZookeeperHosts())) { curator = CuratorFrameworkFactory.builder() .connectString(adapterCanalConfig.getZookeeperHosts()) .retryPolicy(new ExponentialBackoffRetry(1000, 3)) diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java index a7aa1afb..32b39f47 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java @@ -181,11 +181,11 @@ public class CommonRest { */ @PutMapping("/syncSwitch/{destination}/{status}") public Result etl(@PathVariable String destination, @PathVariable String status) { - if ("on".equals(status)) { + if (status.equals("on")) { syncSwitch.on(destination); logger.info("#Destination: {} sync on", destination); return Result.createSuccess("实例: " + destination + " 开启同步成功"); - } else if ("off".equals(status)) { + } else if (status.equals("off")) { syncSwitch.off(destination); logger.info("#Destination: {} sync off", destination); return Result.createSuccess("实例: " + destination + " 关闭同步成功"); diff --git a/client-adapter/pom.xml b/client-adapter/pom.xml index 40006023..7b80e831 100644 --- a/client-adapter/pom.xml +++ b/client-adapter/pom.xml @@ -50,7 +50,7 @@ central - https://repo1.maven.org/maven2 + http://repo1.maven.org/maven2 true @@ -143,7 +143,7 @@ com.h2database h2 - 1.4.196 + 1.4.192 @@ -172,7 +172,7 @@ mysql mysql-connector-java - 5.1.47 + 5.1.48 org.postgresql diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java index 64015268..d27bb544 100644 --- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java +++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/service/RdbSyncService.java @@ -112,12 +112,14 @@ public class RdbSyncService { futures.add(executorThreads[i].submit(() -> { try { - dmlsPartition[j] - .forEach(syncItem -> sync(batchExecutors[j], syncItem.config, syncItem.singleDml)); + dmlsPartition[j].forEach(syncItem -> sync(batchExecutors[j], + syncItem.config, + syncItem.singleDml)); dmlsPartition[j].clear(); batchExecutors[j].commit(); return true; } catch (Throwable e) { + dmlsPartition[j].clear(); batchExecutors[j].rollback(); throw new RuntimeException(e); } @@ -151,50 +153,50 @@ public class RdbSyncService { sync(dmls, dml -> { if (dml.getIsDdl() != null && dml.getIsDdl() && StringUtils.isNotEmpty(dml.getSql())) { // DDL - columnsTypeCache.remove(dml.getDestination() + "." + dml.getDatabase() + "." + dml.getTable()); - return false; + columnsTypeCache.remove(dml.getDestination() + "." + dml.getDatabase() + "." + dml.getTable()); + return false; + } else { + // DML + String destination = StringUtils.trimToEmpty(dml.getDestination()); + String groupId = StringUtils.trimToEmpty(dml.getGroupId()); + String database = dml.getDatabase(); + String table = dml.getTable(); + Map configMap; + if (envProperties != null && !"tcp".equalsIgnoreCase(envProperties.getProperty("canal.conf.mode"))) { + configMap = mappingConfig.get(destination + "-" + groupId + "_" + database + "-" + table); } else { - // DML - String destination = StringUtils.trimToEmpty(dml.getDestination()); - String groupId = StringUtils.trimToEmpty(dml.getGroupId()); - String database = dml.getDatabase(); - String table = dml.getTable(); - Map configMap; - if (envProperties != null && !"tcp".equalsIgnoreCase(envProperties.getProperty("canal.conf.mode"))) { - configMap = mappingConfig.get(destination + "-" + groupId + "_" + database + "-" + table); - } else { - configMap = mappingConfig.get(destination + "_" + database + "-" + table); - } - - if (configMap == null) { - return false; - } - - if (configMap.values().isEmpty()) { - return false; - } - - for (MappingConfig config : configMap.values()) { - boolean caseInsensitive = config.getDbMapping().isCaseInsensitive(); - if (config.getConcurrent()) { - List singleDmls = SingleDml.dml2SingleDmls(dml, caseInsensitive); - singleDmls.forEach(singleDml -> { - int hash = pkHash(config.getDbMapping(), singleDml.getData()); - SyncItem syncItem = new SyncItem(config, singleDml); - dmlsPartition[hash].add(syncItem); - }); - } else { - int hash = 0; - List singleDmls = SingleDml.dml2SingleDmls(dml, caseInsensitive); - singleDmls.forEach(singleDml -> { - SyncItem syncItem = new SyncItem(config, singleDml); - dmlsPartition[hash].add(syncItem); - }); - } - } - return true; + configMap = mappingConfig.get(destination + "_" + database + "-" + table); } - }); + + if (configMap == null) { + return false; + } + + if (configMap.values().isEmpty()) { + return false; + } + + for (MappingConfig config : configMap.values()) { + boolean caseInsensitive = config.getDbMapping().isCaseInsensitive(); + if (config.getConcurrent()) { + List singleDmls = SingleDml.dml2SingleDmls(dml, caseInsensitive); + singleDmls.forEach(singleDml -> { + int hash = pkHash(config.getDbMapping(), singleDml.getData()); + SyncItem syncItem = new SyncItem(config, singleDml); + dmlsPartition[hash].add(syncItem); + }); + } else { + int hash = 0; + List singleDmls = SingleDml.dml2SingleDmls(dml, caseInsensitive); + singleDmls.forEach(singleDml -> { + SyncItem syncItem = new SyncItem(config, singleDml); + dmlsPartition[hash].add(syncItem); + }); + } + } + return true; + } + } ); } /** @@ -208,13 +210,13 @@ public class RdbSyncService { if (config != null) { try { String type = dml.getType(); - if (type != null && "INSERT".equalsIgnoreCase(type)) { + if (type != null && type.equalsIgnoreCase("INSERT")) { insert(batchExecutor, config, dml); - } else if (type != null && "UPDATE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("UPDATE")) { update(batchExecutor, config, dml); - } else if (type != null && "DELETE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("DELETE")) { delete(batchExecutor, config, dml); - } else if (type != null && "TRUNCATE".equalsIgnoreCase(type)) { + } else if (type != null && type.equalsIgnoreCase("TRUNCATE")) { truncate(batchExecutor, config); } if (logger.isDebugEnabled()) { @@ -245,7 +247,10 @@ public class RdbSyncService { StringBuilder insertSql = new StringBuilder(); insertSql.append("INSERT INTO ").append(SyncUtil.getDbTableName(dbMapping)).append(" ("); - columnsMap.forEach((targetColumnName, srcColumnName) -> insertSql.append("`").append(targetColumnName).append("`").append(",")); + columnsMap.forEach((targetColumnName, srcColumnName) -> insertSql.append("`") + .append(targetColumnName) + .append("`") + .append(",")); int len = insertSql.length(); insertSql.delete(len - 1, len).append(") VALUES ("); int mapLen = columnsMap.size(); diff --git a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/support/SyncUtil.java b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/support/SyncUtil.java index 88f1ec94..25247642 100644 --- a/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/support/SyncUtil.java +++ b/client-adapter/rdb/src/main/java/com/alibaba/otter/canal/client/adapter/rdb/support/SyncUtil.java @@ -71,7 +71,7 @@ public class SyncUtil { if (value instanceof Boolean) { pstmt.setBoolean(i, (Boolean) value); } else if (value instanceof String) { - boolean v = !"0".equals(value); + boolean v = !value.equals("0"); pstmt.setBoolean(i, v); } else if (value instanceof Number) { boolean v = ((Number) value).intValue() != 0; diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/CanalConstants.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/CanalConstants.java index 0a9f8256..24c74278 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/CanalConstants.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/CanalConstants.java @@ -12,12 +12,14 @@ public class CanalConstants { public static final String CANAL_FILTER_TRANSACTION_ENTRY = ROOT + "." + "instance.filter.transaction.entry"; - public static final String CANAL_MQ_FLAT_MESSAGE = ROOT + "." + "mq.flat.message"; + public static final String CANAL_MQ_ACCESS_CHANNEL = ROOT + "." + "mq.accessChannel"; + public static final String CANAL_MQ_CANAL_BATCH_SIZE = ROOT + "." + "mq.canalBatchSize"; + public static final String CANAL_MQ_CANAL_GET_TIMEOUT = ROOT + "." + "mq.canalGetTimeout"; + public static final String CANAL_MQ_FLAT_MESSAGE = ROOT + "." + "mq.flatMessage"; + public static final String CANAL_MQ_DATABASE_HASH = ROOT + "." + "mq.database.hash"; - public static final String CANAL_MQ_PARALLEL_THREAD_SIZE = ROOT + "." + "mq.parallel.thread.size"; - public static final String CANAL_MQ_CANAL_BATCH_SIZE = ROOT + "." + "mq.canal.batch.size"; - public static final String CANAL_MQ_CANAL_FETCH_TIMEOUT = ROOT + "." + "mq.canal.fetch.timeout"; - public static final String CANAL_MQ_ACCESS_CHANNEL = ROOT + "." + "mq.access.channel"; + public static final String CANAL_MQ_BUILD_THREAD_SIZE = ROOT + "." + "mq.build.thread.size"; + public static final String CANAL_MQ_SEND_THREAD_SIZE = ROOT + "." + "mq.send.thread.size"; public static final String CANAL_ALIYUN_ACCESS_KEY = ROOT + "." + "aliyun.accessKey"; public static final String CANAL_ALIYUN_SECRET_KEY = ROOT + "." + "aliyun.secretKey"; diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/MQProperties.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/MQProperties.java index fcd7c3db..072d0ebb 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/MQProperties.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/config/MQProperties.java @@ -8,17 +8,18 @@ package com.alibaba.otter.canal.connector.core.config; */ public class MQProperties { - private boolean flatMessage = true; - private boolean databaseHash = true; - private boolean filterTransactionEntry = true; - private Integer parallelThreadSize = 8; - private Integer fetchTimeout = 100; - private Integer batchSize = 50; - private String accessChannel = "local"; + private boolean flatMessage = true; + private boolean databaseHash = true; + private boolean filterTransactionEntry = true; + private Integer parallelBuildThreadSize = 8; + private Integer parallelSendThreadSize = 30; + private Integer fetchTimeout = 100; + private Integer batchSize = 50; + private String accessChannel = "local"; - private String aliyunAccessKey = ""; - private String aliyunSecretKey = ""; - private int aliyunUid = 0; + private String aliyunAccessKey = ""; + private String aliyunSecretKey = ""; + private int aliyunUid = 0; public boolean isFlatMessage() { return flatMessage; @@ -44,12 +45,20 @@ public class MQProperties { this.filterTransactionEntry = filterTransactionEntry; } - public Integer getParallelThreadSize() { - return parallelThreadSize; + public Integer getParallelBuildThreadSize() { + return parallelBuildThreadSize; } - public void setParallelThreadSize(Integer parallelThreadSize) { - this.parallelThreadSize = parallelThreadSize; + public void setParallelBuildThreadSize(Integer parallelBuildThreadSize) { + this.parallelBuildThreadSize = parallelBuildThreadSize; + } + + public Integer getParallelSendThreadSize() { + return parallelSendThreadSize; + } + + public void setParallelSendThreadSize(Integer parallelSendThreadSize) { + this.parallelSendThreadSize = parallelSendThreadSize; } public Integer getFetchTimeout() { diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/AbstractMQProducer.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/AbstractMQProducer.java index 98e35381..940e8fc5 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/AbstractMQProducer.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/AbstractMQProducer.java @@ -22,20 +22,30 @@ public abstract class AbstractMQProducer implements CanalMQProducer { protected MQProperties mqProperties; - protected ThreadPoolExecutor executor; + protected ThreadPoolExecutor sendExecutor; + protected ThreadPoolExecutor buildExecutor; @Override public void init(Properties properties) { // parse canal mq properties loadCanalMqProperties(properties); - int parallelThreadSize = mqProperties.getParallelThreadSize(); - executor = new ThreadPoolExecutor(parallelThreadSize, - parallelThreadSize, + int parallelBuildThreadSize = mqProperties.getParallelBuildThreadSize(); + buildExecutor = new ThreadPoolExecutor(parallelBuildThreadSize, + parallelBuildThreadSize, 0, TimeUnit.SECONDS, - new ArrayBlockingQueue(parallelThreadSize * 2), - new NamedThreadFactory("MQParallel"), + new ArrayBlockingQueue(parallelBuildThreadSize * 2), + new NamedThreadFactory("MQ-Parallel-Builder"), + new ThreadPoolExecutor.CallerRunsPolicy()); + + int parallelSendThreadSize = mqProperties.getParallelSendThreadSize(); + sendExecutor = new ThreadPoolExecutor(parallelSendThreadSize, + parallelSendThreadSize, + 0, + TimeUnit.SECONDS, + new ArrayBlockingQueue(parallelSendThreadSize * 2), + new NamedThreadFactory("MQ-Parallel-Sender"), new ThreadPoolExecutor.CallerRunsPolicy()); } @@ -46,8 +56,12 @@ public abstract class AbstractMQProducer implements CanalMQProducer { @Override public void stop() { - if (executor != null) { - executor.shutdownNow(); + if (buildExecutor != null) { + buildExecutor.shutdownNow(); + } + + if (sendExecutor != null) { + sendExecutor.shutdownNow(); } } @@ -57,7 +71,8 @@ public abstract class AbstractMQProducer implements CanalMQProducer { * canal.mq.flat.message = true * canal.mq.database.hash = true * canal.mq.filter.transaction.entry = true - * canal.mq.parallel.thread.size = 8 + * canal.mq.parallel.build.thread.size = 8 + * canal.mq.parallel.send.thread.size = 8 * canal.mq.batch.size = 50 * canal.mq.timeout = 100 * canal.mq.access.channel = local @@ -70,6 +85,7 @@ public abstract class AbstractMQProducer implements CanalMQProducer { if (!StringUtils.isEmpty(flatMessage)) { mqProperties.setFlatMessage(Boolean.parseBoolean(flatMessage)); } + String databaseHash = properties.getProperty(CanalConstants.CANAL_MQ_DATABASE_HASH); if (!StringUtils.isEmpty(databaseHash)) { mqProperties.setDatabaseHash(Boolean.parseBoolean(databaseHash)); @@ -78,15 +94,19 @@ public abstract class AbstractMQProducer implements CanalMQProducer { if (!StringUtils.isEmpty(filterTranEntry)) { mqProperties.setFilterTransactionEntry(Boolean.parseBoolean(filterTranEntry)); } - String parallelThreadSize = properties.getProperty(CanalConstants.CANAL_MQ_PARALLEL_THREAD_SIZE); - if (!StringUtils.isEmpty(parallelThreadSize)) { - mqProperties.setParallelThreadSize(Integer.parseInt(parallelThreadSize)); + String parallelBuildThreadSize = properties.getProperty(CanalConstants.CANAL_MQ_BUILD_THREAD_SIZE); + if (!StringUtils.isEmpty(parallelBuildThreadSize)) { + mqProperties.setParallelBuildThreadSize(Integer.parseInt(parallelBuildThreadSize)); + } + String parallelSendThreadSize = properties.getProperty(CanalConstants.CANAL_MQ_SEND_THREAD_SIZE); + if (!StringUtils.isEmpty(parallelSendThreadSize)) { + mqProperties.setParallelSendThreadSize(Integer.parseInt(parallelSendThreadSize)); } String batchSize = properties.getProperty(CanalConstants.CANAL_MQ_CANAL_BATCH_SIZE); if (!StringUtils.isEmpty(batchSize)) { mqProperties.setBatchSize(Integer.parseInt(batchSize)); } - String timeOut = properties.getProperty(CanalConstants.CANAL_MQ_CANAL_FETCH_TIMEOUT); + String timeOut = properties.getProperty(CanalConstants.CANAL_MQ_CANAL_GET_TIMEOUT); if (!StringUtils.isEmpty(timeOut)) { mqProperties.setFetchTimeout(Integer.parseInt(timeOut)); } @@ -107,4 +127,14 @@ public abstract class AbstractMQProducer implements CanalMQProducer { mqProperties.setAliyunUid(Integer.parseInt(aliyunUid)); } } + + /** + * 兼容下<=1.1.4的mq配置项 + */ + protected void doMoreCompatibleConvert(String oldKey, String newKey, Properties properties) { + String value = properties.getProperty(oldKey); + if (StringUtils.isNotEmpty(value)) { + properties.setProperty(newKey, value); + } + } } diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java index e1678698..ad2b4c87 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/producer/MQMessageUtils.java @@ -49,7 +49,7 @@ public class MQMessageUtils { int i = pkHashConfig.lastIndexOf(":"); if (i > 0) { String pkStr = pkHashConfig.substring(i + 1); - if ("$pk$".equalsIgnoreCase(pkStr)) { + if (pkStr.equalsIgnoreCase("$pk$")) { data.hashMode.autoPkHash = true; } else { data.hashMode.pkNames = Lists.newArrayList(StringUtils.split(pkStr, diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/JdbcTypeUtil.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/JdbcTypeUtil.java index ffca31a6..50053f49 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/JdbcTypeUtil.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/JdbcTypeUtil.java @@ -78,7 +78,7 @@ public class JdbcTypeUtil { public static Object typeConvert(String tableName, String columnName, String value, int sqlType, String mysqlType) { if (value == null - || ("".equals(value) && !(isText(mysqlType) || sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.LONGVARCHAR))) { + || (value.equals("") && !(isText(mysqlType) || sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.LONGVARCHAR))) { return null; } diff --git a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/MessageUtil.java b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/MessageUtil.java index 160f04e3..1096739b 100644 --- a/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/MessageUtil.java +++ b/connector/core/src/main/java/com/alibaba/otter/canal/connector/core/util/MessageUtil.java @@ -42,6 +42,7 @@ public class MessageUtil { CanalEntry.EventType eventType = rowChange.getEventType(); final CommonMessage msg = new CommonMessage(); + msg.setIsDdl(rowChange.getIsDdl()); msg.setDatabase(entry.getHeader().getSchemaName()); msg.setTable(entry.getHeader().getTableName()); msg.setType(eventType.toString()); diff --git a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/config/KafkaConstants.java b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/config/KafkaConstants.java index 8a6a1237..cbfc2a4d 100644 --- a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/config/KafkaConstants.java +++ b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/config/KafkaConstants.java @@ -8,9 +8,9 @@ package com.alibaba.otter.canal.connector.kafka.config; */ public class KafkaConstants { - public static final String ROOT = "canal"; + public static final String ROOT = "kafka"; - public static final String CANAL_MQ_KAFKA_KERBEROS_ENABLE = ROOT + "." + "mq.kafka.kerberos.enable"; - public static final String CANAL_MQ_KAFKA_KERBEROS_KRB5_FILE = ROOT + "." + "mq.kafka.kerberos.krb5.file"; - public static final String CANAL_MQ_KAFKA_KERBEROS_JAAS_FILE = ROOT + "." + "mq.kafka.kerberos.jaas.file"; + public static final String CANAL_MQ_KAFKA_KERBEROS_ENABLE = ROOT + "." + "kerberos.enable"; + public static final String CANAL_MQ_KAFKA_KERBEROS_KRB5_FILE = ROOT + "." + "kerberos.krb5.file"; + public static final String CANAL_MQ_KAFKA_KERBEROS_JAAS_FILE = ROOT + "." + "kerberos.jaas.file"; } diff --git a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java index 8f3a70c7..df50bc28 100644 --- a/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java +++ b/connector/kafka-connector/src/main/java/com/alibaba/otter/canal/connector/kafka/producer/CanalKafkaProducer.java @@ -59,6 +59,7 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro Properties kafkaProperties = new Properties(); kafkaProperties.putAll(kafkaProducerConfig.getKafkaProperties()); + kafkaProperties.put("max.in.flight.requests.per.connection", 1); kafkaProperties.put("key.serializer", StringSerializer.class); if (kafkaProducerConfig.isKerberosEnabled()) { File krb5File = new File(kafkaProducerConfig.getKrb5File()); @@ -83,6 +84,19 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro private void loadKafkaProperties(Properties properties) { KafkaProducerConfig kafkaProducerConfig = (KafkaProducerConfig) this.mqProperties; Map kafkaProperties = kafkaProducerConfig.getKafkaProperties(); + // 兼容下<=1.1.4的mq配置 + doMoreCompatibleConvert("canal.mq.servers", "kafka.bootstrap.servers", properties); + doMoreCompatibleConvert("canal.mq.acks", "kafka.acks", properties); + doMoreCompatibleConvert("canal.mq.compressionType", "kafka.compression.type", properties); + doMoreCompatibleConvert("canal.mq.retries", "kafka.retries", properties); + doMoreCompatibleConvert("canal.mq.batchSize", "kafka.batch.size", properties); + doMoreCompatibleConvert("canal.mq.lingerMs", "kafka.linger.ms", properties); + doMoreCompatibleConvert("canal.mq.maxRequestSize", "kafka.max.request.size", properties); + doMoreCompatibleConvert("canal.mq.bufferMemory", "kafka.buffer.memory", properties); + doMoreCompatibleConvert("canal.mq.kafka.kerberos.enable", "kafka.kerberos.enable", properties); + doMoreCompatibleConvert("canal.mq.kafka.kerberos.krb5.file", "kafka.kerberos.krb5.file", properties); + doMoreCompatibleConvert("canal.mq.kafka.kerberos.jaas.file", "kafka.kerberos.jaas.file", properties); + for (Map.Entry entry : properties.entrySet()) { String key = (String) entry.getKey(); Object value = entry.getValue(); @@ -91,7 +105,6 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro kafkaProperties.put(key, value); } } - String kerberosEnabled = properties.getProperty(KafkaConstants.CANAL_MQ_KAFKA_KERBEROS_ENABLE); if (!StringUtils.isEmpty(kerberosEnabled)) { kafkaProducerConfig.setKerberosEnabled(Boolean.parseBoolean(kerberosEnabled)); @@ -123,7 +136,7 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro @Override public void send(MQDestination mqDestination, Message message, Callback callback) { - ExecutorTemplate template = new ExecutorTemplate(executor); + ExecutorTemplate template = new ExecutorTemplate(sendExecutor); try { List result; @@ -186,7 +199,7 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro if (!flat) { if (mqDestination.getPartitionHash() != null && !mqDestination.getPartitionHash().isEmpty()) { // 并发构造 - EntryRowData[] datas = MQMessageUtils.buildMessageData(message, executor); + EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); // 串行分区 Message[] messages = MQMessageUtils.messagePartition(datas, message.getId(), @@ -200,7 +213,8 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro records.add(new ProducerRecord<>(topicName, i, null, - CanalMessageSerializerUtil.serializer(messagePartition, true))); + CanalMessageSerializerUtil.serializer(messagePartition, + mqProperties.isFilterTransactionEntry()))); } } } else { @@ -208,12 +222,12 @@ public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQPro records.add(new ProducerRecord<>(topicName, partition, null, - CanalMessageSerializerUtil.serializer(message, true))); + CanalMessageSerializerUtil.serializer(message, mqProperties.isFilterTransactionEntry()))); } } else { // 发送扁平数据json // 并发构造 - EntryRowData[] datas = MQMessageUtils.buildMessageData(message, executor); + EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); // 串行分区 List flatMessages = MQMessageUtils.messageConverter(datas, message.getId()); for (FlatMessage flatMessage : flatMessages) { diff --git a/connector/pom.xml b/connector/pom.xml index 46b4a6b5..b56937d5 100644 --- a/connector/pom.xml +++ b/connector/pom.xml @@ -41,7 +41,7 @@ central - https://repo1.maven.org/maven2 + http://repo1.maven.org/maven2 true diff --git a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java index 72d9ffc6..c7561fe6 100644 --- a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java +++ b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java @@ -201,13 +201,6 @@ public class CanalRabbitMQConsumer implements CanalMsgConsumer { @Override public void disconnect() { - if (connect != null) { - try { - connect.close(); - } catch (IOException e) { - throw new CanalClientException("stop connect error", e); - } - } if (channel != null) { try { channel.close(); @@ -215,5 +208,13 @@ public class CanalRabbitMQConsumer implements CanalMsgConsumer { throw new CanalClientException("stop channel error", e); } } + + if (connect != null) { + try { + connect.close(); + } catch (IOException e) { + throw new CanalClientException("stop connect error", e); + } + } } } diff --git a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/producer/CanalRabbitMQProducer.java b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/producer/CanalRabbitMQProducer.java index 1374dea2..36e795e4 100644 --- a/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/producer/CanalRabbitMQProducer.java +++ b/connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/producer/CanalRabbitMQProducer.java @@ -81,6 +81,8 @@ public class CanalRabbitMQProducer extends AbstractMQProducer implements CanalMQ private void loadRabbitMQProperties(Properties properties) { RabbitMQProducerConfig rabbitMQProperties = (RabbitMQProducerConfig) this.mqProperties; + // 兼容下<=1.1.4的mq配置 + doMoreCompatibleConvert("canal.mq.servers", "rabbitmq.host", properties); String host = properties.getProperty(RabbitMQConstants.RABBITMQ_HOST); if (!StringUtils.isEmpty(host)) { @@ -106,7 +108,7 @@ public class CanalRabbitMQProducer extends AbstractMQProducer implements CanalMQ @Override public void send(final MQDestination destination, Message message, Callback callback) { - ExecutorTemplate template = new ExecutorTemplate(executor); + ExecutorTemplate template = new ExecutorTemplate(sendExecutor); try { if (!StringUtils.isEmpty(destination.getDynamicTopic())) { // 动态topic @@ -149,7 +151,7 @@ public class CanalRabbitMQProducer extends AbstractMQProducer implements CanalMQ sendMessage(topicName, message); } else { // 并发构造 - MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(messageSub, executor); + MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(messageSub, buildExecutor); // 串行分区 List flatMessages = MQMessageUtils.messageConverter(datas, messageSub.getId()); for (FlatMessage flatMessage : flatMessages) { diff --git a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java index 08a62edf..3793ca73 100644 --- a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java +++ b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java @@ -47,7 +47,6 @@ public class CanalRocketMQConsumer implements CanalMsgConsumer { private String nameServer; private String topic; private String groupName; - private volatile boolean connected = false; private DefaultMQPushConsumer rocketMQConsumer; private BlockingQueue> messageBlockingQueue; private int batchSize = -1; @@ -135,7 +134,6 @@ public class CanalRocketMQConsumer implements CanalMsgConsumer { }); rocketMQConsumer.start(); } catch (MQClientException ex) { - connected = false; logger.error("Start RocketMQ consumer error", ex); } } @@ -231,6 +229,5 @@ public class CanalRocketMQConsumer implements CanalMsgConsumer { public void disconnect() { rocketMQConsumer.unsubscribe(topic); rocketMQConsumer.shutdown(); - connected = false; } } diff --git a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/producer/CanalRocketMQProducer.java b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/producer/CanalRocketMQProducer.java index da293f4b..fd2fd304 100644 --- a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/producer/CanalRocketMQProducer.java +++ b/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/producer/CanalRocketMQProducer.java @@ -88,6 +88,11 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ private void loadRocketMQProperties(Properties properties) { RocketMQProducerConfig rocketMQProperties = (RocketMQProducerConfig) this.mqProperties; + // 兼容下<=1.1.4的mq配置 + doMoreCompatibleConvert("canal.mq.servers", "rocketmq.namesrv.addr", properties); + doMoreCompatibleConvert("canal.mq.producerGroup", "rocketmq.producer.group", properties); + doMoreCompatibleConvert("canal.mq.namespace", "rocketmq.namespace", properties); + doMoreCompatibleConvert("canal.mq.retries", "rocketmq.retry.times.when.send.failed", properties); String producerGroup = properties.getProperty(RocketMQConstants.ROCKETMQ_PRODUCER_GROUP); if (!StringUtils.isEmpty(producerGroup)) { @@ -121,7 +126,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ @Override public void send(MQDestination destination, com.alibaba.otter.canal.protocol.Message message, Callback callback) { - ExecutorTemplate template = new ExecutorTemplate(executor); + ExecutorTemplate template = new ExecutorTemplate(sendExecutor); try { if (!StringUtils.isEmpty(destination.getDynamicTopic())) { // 动态topic @@ -159,7 +164,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ if (!mqProperties.isFlatMessage()) { if (destination.getPartitionHash() != null && !destination.getPartitionHash().isEmpty()) { // 并发构造 - MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(message, executor); + MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); // 串行分区 com.alibaba.otter.canal.protocol.Message[] messages = MQMessageUtils.messagePartition(datas, message.getId(), @@ -168,7 +173,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ mqProperties.isDatabaseHash()); int length = messages.length; - ExecutorTemplate template = new ExecutorTemplate(executor); + ExecutorTemplate template = new ExecutorTemplate(sendExecutor); for (int i = 0; i < length; i++) { com.alibaba.otter.canal.protocol.Message dataPartition = messages[i]; if (dataPartition != null) { @@ -190,7 +195,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ } } else { // 并发构造 - MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(message, executor); + MQMessageUtils.EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); // 串行分区 List flatMessages = MQMessageUtils.messageConverter(datas, message.getId()); // 初始化分区合并队列 @@ -211,7 +216,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ } } - ExecutorTemplate template = new ExecutorTemplate(executor); + ExecutorTemplate template = new ExecutorTemplate(sendExecutor); for (int i = 0; i < partitionFlatMessages.size(); i++) { final List flatMessagePart = partitionFlatMessages.get(i); if (flatMessagePart != null) { @@ -244,7 +249,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ private void sendMessage(Message message, int partition) { try { SendResult sendResult = this.defaultMQProducer.send(message, (mqs, msg, arg) -> { - if (partition > mqs.size()) { + if (partition >= mqs.size()) { return mqs.get(partition % mqs.size()); } else { return mqs.get(partition); @@ -283,7 +288,7 @@ public class CanalRocketMQProducer extends AbstractMQProducer implements CanalMQ } } else { MessageQueue queue; - if (partition > size) { + if (partition >= size) { queue = queues.get(partition % size); } else { queue = queues.get(partition); diff --git a/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/TableMapLogEvent.java b/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/TableMapLogEvent.java index 18201ea9..79fa4147 100644 --- a/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/TableMapLogEvent.java +++ b/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/TableMapLogEvent.java @@ -326,8 +326,8 @@ public final class TableMapLogEvent extends LogEvent { * * Source : http://forge.mysql.com/wiki/MySQL_Internals_Binary_Log */ - protected final String dbname; - protected final String tblname; + protected String dbname; + protected String tblname; /** * Holding mysql column information. @@ -418,6 +418,7 @@ public final class TableMapLogEvent extends LogEvent { buffer.position(commonHeaderLen + postHeaderLen); dbname = buffer.getString(); buffer.forward(1); /* termination null */ + // fixed issue #2714 tblname = buffer.getString(); buffer.forward(1); /* termination null */ @@ -792,6 +793,22 @@ public final class TableMapLogEvent extends LogEvent { return tblname; } + public String getDbname() { + return dbname; + } + + public void setDbname(String dbname) { + this.dbname = dbname; + } + + public String getTblname() { + return tblname; + } + + public void setTblname(String tblname) { + this.tblname = tblname; + } + public final int getColumnCnt() { return columnCnt; } diff --git a/deployer/src/main/bin/restart.sh b/deployer/src/main/bin/restart.sh index 7b0ed7b4..e2499d51 100644 --- a/deployer/src/main/bin/restart.sh +++ b/deployer/src/main/bin/restart.sh @@ -2,5 +2,5 @@ args=$@ -bash stop.sh $args -bash startup.sh $args +sh stop.sh $args +sh startup.sh $args \ No newline at end of file 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 7f129831..0a0a3980 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 @@ -443,11 +443,11 @@ public class CanalController { } config.setManagerAddress(managerAddress); } - } else if (config.getMode().isSpring()) { - String springXml = getProperty(properties, CanalConstants.getInstancSpringXmlKey(destination)); - if (StringUtils.isNotEmpty(springXml)) { - config.setSpringXml(springXml); - } + } + + String springXml = getProperty(properties, CanalConstants.getInstancSpringXmlKey(destination)); + if (StringUtils.isNotEmpty(springXml)) { + config.setSpringXml(springXml); } return config; diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/admin/CanalAdminController.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/admin/CanalAdminController.java index 5877f7fe..5144e7f7 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/admin/CanalAdminController.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/admin/CanalAdminController.java @@ -218,11 +218,6 @@ public class CanalAdminController implements CanalAdmin { return FileUtils.readFileFromOffset("../logs/" + destination + "/" + fileName, lines, "UTF-8"); } - @Override - public String instanceMeta(String destination, String fileName) { - return FileUtils.readFileFromOffset("../conf/" + destination + "/" + fileName, 100, "UTF-8"); - } - private InstanceAction getInstanceAction(String destination) { Map monitors = canalStater.getController() .getInstanceConfigMonitors(); diff --git a/deployer/src/main/resources/canal.properties b/deployer/src/main/resources/canal.properties index 234ffb00..a14766cd 100644 --- a/deployer/src/main/resources/canal.properties +++ b/deployer/src/main/resources/canal.properties @@ -109,19 +109,21 @@ canal.instance.global.spring.xml = classpath:spring/file-instance.xml ################################################## ######### MQ Properties ############# ################################################## -canal.mq.flat.message = true -canal.mq.database.hash = true -canal.mq.parallel.thread.size = 8 -canal.mq.canal.batch.size = 50 -canal.mq.canal.fetch.timeout = 100 -# Set this value to "cloud", if you want open message trace feature in aliyun. -canal.mq.access.channel = local - # aliyun ak/sk , support rds/mq canal.aliyun.accessKey = canal.aliyun.secretKey = canal.aliyun.uid= +canal.mq.flatMessage = true +canal.mq.canalBatchSize = 50 +canal.mq.canalGetTimeout = 100 +# Set this value to "cloud", if you want open message trace feature in aliyun. +canal.mq.accessChannel = local + +canal.mq.database.hash = true +canal.mq.send.thread.size = 30 +canal.mq.build.thread.size = 8 + ################################################## ######### Kafka ############# ################################################## @@ -135,9 +137,9 @@ kafka.buffer.memory = 33554432 kafka.max.in.flight.requests.per.connection = 1 kafka.retries = 0 -canal.mq.kafka.kerberos.enable = false -canal.mq.kafka.kerberos.krb5.file = "../conf/kerberos/krb5.conf" -canal.mq.kafka.kerberos.jaas.file = "../conf/kerberos/jaas.conf" +kafka.kerberos.enable = false +kafka.kerberos.krb5.file = "../conf/kerberos/krb5.conf" +kafka.kerberos.jaas.file = "../conf/kerberos/jaas.conf" ################################################## ######### RocketMQ ############# diff --git a/deployer/src/main/resources/example/instance.properties b/deployer/src/main/resources/example/instance.properties index 5ceccab6..06e3dad3 100644 --- a/deployer/src/main/resources/example/instance.properties +++ b/deployer/src/main/resources/example/instance.properties @@ -40,7 +40,7 @@ canal.instance.enableDruid=false # table regex canal.instance.filter.regex=.*\\..* # table black regex -canal.instance.filter.black.regex= +canal.instance.filter.black.regex=mysql\\.slave_.* # table field filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2) #canal.instance.filter.field=test1.t_product:id/subject/keywords,test2.t_company:id/name/contact/ch # table field black filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2) diff --git a/example/src/main/java/com/alibaba/otter/canal/example/AbstractCanalClientTest.java b/example/src/main/java/com/alibaba/otter/canal/example/AbstractCanalClientTest.java index 1f77d00f..25746a26 100644 --- a/example/src/main/java/com/alibaba/otter/canal/example/AbstractCanalClientTest.java +++ b/example/src/main/java/com/alibaba/otter/canal/example/AbstractCanalClientTest.java @@ -76,16 +76,17 @@ public class AbstractCanalClientTest extends BaseCanalClientTest { if (batchId != -1) { connector.ack(batchId); // 提交确认 - // connector.rollback(batchId); // 处理失败, 回滚数据 } } - } catch (Exception e) { + } catch (Throwable e) { logger.error("process error!", e); try { Thread.sleep(1000L); } catch (InterruptedException e1) { // ignore } + + connector.rollback(); // 处理失败, 回滚数据 } finally { connector.disconnect(); MDC.remove("destination"); diff --git a/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientFlatMessageExample.java b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientFlatMessageExample.java index 5d5fd2d1..60b798ac 100644 --- a/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientFlatMessageExample.java +++ b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientFlatMessageExample.java @@ -117,8 +117,7 @@ public class CanalKafkaClientFlatMessageExample { } for (FlatMessage message : messages) { long batchId = message.getId(); - int size = message.getData() == null ? 0 : message.getData().size(); - if (batchId == -1 || size == 0) { + if (batchId == -1 || message.getData() == null) { // try { // Thread.sleep(1000); // } catch (InterruptedException e) { 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 9752645a..8a35326d 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 @@ -5,7 +5,6 @@ import java.io.FilenameFilter; import java.net.InetSocketAddress; import java.net.URL; import java.net.URLClassLoader; -import java.nio.charset.Charset; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -294,7 +293,7 @@ public class CanalInstanceWithManager extends AbstractCanalInstance { } mysqlEventParser.setDestination(destination); // 编码参数 - mysqlEventParser.setConnectionCharset(Charset.forName(parameters.getConnectionCharset())); + mysqlEventParser.setConnectionCharset(parameters.getConnectionCharset()); mysqlEventParser.setConnectionCharsetNumber(parameters.getConnectionCharsetNumber()); // 网络相关参数 mysqlEventParser.setDefaultConnectionTimeoutInSeconds(parameters.getDefaultConnectionTimeoutInSeconds()); @@ -375,7 +374,7 @@ public class CanalInstanceWithManager extends AbstractCanalInstance { LocalBinlogEventParser localBinlogEventParser = new LocalBinlogEventParser(); localBinlogEventParser.setDestination(destination); localBinlogEventParser.setBufferSize(parameters.getReceiveBufferSize()); - localBinlogEventParser.setConnectionCharset(Charset.forName(parameters.getConnectionCharset())); + localBinlogEventParser.setConnectionCharset(parameters.getConnectionCharset()); localBinlogEventParser.setConnectionCharsetNumber(parameters.getConnectionCharsetNumber()); localBinlogEventParser.setDirectory(parameters.getLocalBinlogDirectory()); localBinlogEventParser.setProfilingEnabled(false); 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 49543939..28826989 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 @@ -130,9 +130,9 @@ public class FileMixedMetaManager extends MemoryMetaManager implements CanalMeta } public void stop() { - super.stop(); - flushDataToFile();// 刷新数据 + + super.stop(); executor.shutdownNow(); destinations.clear(); batches.clear(); diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java index fbde7790..7ae01f63 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java @@ -50,7 +50,7 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser { if (eventBlackFilter != null && eventBlackFilter instanceof AviaterRegexFilter) { convert.setNameBlackFilter((AviaterRegexFilter) eventBlackFilter); } - + convert.setFieldFilterMap(getFieldFilterMap()); convert.setFieldBlackFilterMap(getFieldBlackFilterMap()); @@ -93,13 +93,13 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser { } } } - + @Override public void setFieldFilter(String fieldFilter) { - super.setFieldFilter(fieldFilter); - - // 触发一下filter变更 - if (binlogParser instanceof LogEventConvert) { + super.setFieldFilter(fieldFilter); + + // 触发一下filter变更 + if (binlogParser instanceof LogEventConvert) { ((LogEventConvert) binlogParser).setFieldFilterMap(getFieldFilterMap()); } @@ -107,13 +107,13 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser { ((DatabaseTableMeta) tableMetaTSDB).setFieldFilterMap(getFieldFilterMap()); } } - + @Override public void setFieldBlackFilter(String fieldBlackFilter) { - super.setFieldBlackFilter(fieldBlackFilter); - - // 触发一下filter变更 - if (binlogParser instanceof LogEventConvert) { + super.setFieldBlackFilter(fieldBlackFilter); + + // 触发一下filter变更 + if (binlogParser instanceof LogEventConvert) { ((LogEventConvert) binlogParser).setFieldBlackFilterMap(getFieldBlackFilterMap()); } @@ -184,11 +184,15 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser { this.connectionCharsetNumber = connectionCharsetNumber; } - public void setConnectionCharset(Charset connectionCharset) { + public void setConnectionCharsetStd(Charset connectionCharset) { this.connectionCharset = connectionCharset; } public void setConnectionCharset(String connectionCharset) { + if ("UTF8MB4".equalsIgnoreCase(connectionCharset)) { + connectionCharset = "UTF-8"; + } + this.connectionCharset = Charset.forName(connectionCharset); } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlConnection.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlConnection.java index 785f2950..59dea37a 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlConnection.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlConnection.java @@ -393,13 +393,13 @@ public class MysqlConnection implements ErosaConnection { logger.warn("update wait_timeout failed", e); } try { - update("set net_write_timeout=1800"); + update("set net_write_timeout=7200"); } catch (Exception e) { logger.warn("update net_write_timeout failed", e); } try { - update("set net_read_timeout=1800"); + update("set net_read_timeout=7200"); } catch (Exception e) { logger.warn("update net_read_timeout failed", e); } @@ -522,7 +522,7 @@ public class MysqlConnection implements ErosaConnection { rs = query("select @@global.binlog_checksum"); List columnValues = rs.getFieldValues(); if (columnValues != null && columnValues.size() >= 1 && columnValues.get(0) != null - && "CRC32".equals(columnValues.get(0).toUpperCase())) { + && columnValues.get(0).toUpperCase().equals("CRC32")) { binlogChecksum = LogEvent.BINLOG_CHECKSUM_ALG_CRC32; } else { binlogChecksum = LogEvent.BINLOG_CHECKSUM_ALG_OFF; diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java index cf8d3d4f..69bd75ca 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java @@ -20,7 +20,6 @@ import org.slf4j.LoggerFactory; import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter; import com.alibaba.otter.canal.parse.exception.CanalParseException; -import com.taobao.tddl.dbsync.binlog.exception.TableIdNotFoundException; import com.alibaba.otter.canal.parse.inbound.BinlogParser; import com.alibaba.otter.canal.parse.inbound.TableMeta; import com.alibaba.otter.canal.parse.inbound.TableMeta.FieldMeta; @@ -59,6 +58,7 @@ import com.taobao.tddl.dbsync.binlog.event.UserVarLogEvent; import com.taobao.tddl.dbsync.binlog.event.WriteRowsLogEvent; import com.taobao.tddl.dbsync.binlog.event.XidLogEvent; import com.taobao.tddl.dbsync.binlog.event.mariadb.AnnotateRowsEvent; +import com.taobao.tddl.dbsync.binlog.exception.TableIdNotFoundException; /** * 基于{@linkplain LogEvent}转化为Entry对象的处理 @@ -88,8 +88,8 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar private volatile AviaterRegexFilter nameFilter; // 运行时引用可能会有变化,比如规则发生变化时 private volatile AviaterRegexFilter nameBlackFilter; - private Map> fieldFilterMap = new HashMap>(); - private Map> fieldBlackFilterMap = new HashMap>(); + private Map> fieldFilterMap = new HashMap>(); + private Map> fieldBlackFilterMap = new HashMap>(); private TableMetaCache tableMetaCache; private Charset charset = Charset.defaultCharset(); @@ -119,6 +119,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar case LogEvent.XID_EVENT: return parseXidEvent((XidLogEvent) logEvent); case LogEvent.TABLE_MAP_EVENT: + parseTableMapEvent((TableMapLogEvent) logEvent); break; case LogEvent.WRITE_ROWS_EVENT_V1: case LogEvent.WRITE_ROWS_EVENT: @@ -144,6 +145,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar return parseGTIDLogEvent((GtidLogEvent) logEvent); case LogEvent.HEARTBEAT_LOG_EVENT: return parseHeartbeatLogEvent((HeartbeatLogEvent) logEvent); + default: break; } @@ -275,7 +277,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar Header header = createHeader(event.getHeader(), schemaName, tableName, type); RowChange.Builder rowChangeBuider = RowChange.newBuilder(); - if (type != EventType.QUERY && !isDml) { + if (type == EventType.QUERY && !isDml) { rowChangeBuider.setIsDdl(true); } rowChangeBuider.setSql(queryString); @@ -491,6 +493,18 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar return parseRowsEvent(event, null); } + public void parseTableMapEvent(TableMapLogEvent event) { + try { + String charsetDbName = new String(event.getDbName().getBytes(ISO_8859_1), charset.name()); + event.setDbname(charsetDbName); + + String charsetTbName = new String(event.getTableName().getBytes(ISO_8859_1), charset.name()); + event.setTblname(charsetTbName); + } catch (UnsupportedEncodingException e) { + throw new CanalParseException(e); + } + } + public Entry parseRowsEvent(RowsLogEvent event, TableMeta tableMeta) { if (filterRows) { return null; @@ -588,15 +602,15 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar boolean tableError = false; // check table fileds count,只能处理加字段 boolean existRDSNoPrimaryKey = false; - //获取字段过滤条件 + // 获取字段过滤条件 List fieldList = null; List blackFieldList = null; - + if (tableMeta != null) { - fieldList = fieldFilterMap.get(tableMeta.getFullName().toUpperCase()); - blackFieldList = fieldBlackFilterMap.get(tableMeta.getFullName().toUpperCase()); + fieldList = fieldFilterMap.get(tableMeta.getFullName().toUpperCase()); + blackFieldList = fieldBlackFilterMap.get(tableMeta.getFullName().toUpperCase()); } - + if (tableMeta != null && columnInfo.length > tableMeta.getFields().size()) { if (tableMetaCache.isOnRDS()) { // 特殊处理下RDS的场景 @@ -650,7 +664,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar if (existRDSNoPrimaryKey && i == columnCnt - 1 && info.type == LogEvent.MYSQL_TYPE_LONGLONG) { // 不解析最后一列 - String rdsRowIdColumnName = "#alibaba_rds_row_id#"; + String rdsRowIdColumnName = "__#alibaba_rds_row_id#__"; buffer.nextValue(rdsRowIdColumnName, i, info.type, info.meta, false); Column.Builder columnBuilder = Column.newBuilder(); columnBuilder.setName(rdsRowIdColumnName); @@ -664,7 +678,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar columnBuilder.setUpdated(false); if (needField(fieldList, blackFieldList, columnBuilder.getName())) { - if (isAfter) { + if (isAfter) { rowDataBuilder.addAfterColumns(columnBuilder.build()); } else { rowDataBuilder.addBeforeColumns(columnBuilder.build()); @@ -712,7 +726,6 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar // fixed issue // https://github.com/alibaba/canal/issues/66,特殊处理binary/varbinary,不能做编码处理 boolean isBinary = false; - boolean isSingleBit = false; if (fieldMeta != null) { if (StringUtils.containsIgnoreCase(fieldMeta.getColumnType(), "VARBINARY")) { isBinary = true; @@ -830,7 +843,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar columnBuilder.getIsNull() ? null : columnBuilder.getValue(), i)); if (needField(fieldList, blackFieldList, columnBuilder.getName())) { - if (isAfter) { + if (isAfter) { rowDataBuilder.addAfterColumns(columnBuilder.build()); } else { rowDataBuilder.addBeforeColumns(columnBuilder.build()); @@ -969,16 +982,17 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar private boolean isRDSHeartBeat(String schema, String table) { return "mysql".equalsIgnoreCase(schema) && "ha_health_check".equalsIgnoreCase(table); } - + /** * 字段过滤判断 */ private boolean needField(List fieldList, List blackFieldList, String columnName) { - if (fieldList == null || fieldList.isEmpty()) { - return blackFieldList == null || blackFieldList.isEmpty() || !blackFieldList.contains(columnName.toUpperCase()); - } else { - return fieldList.contains(columnName.toUpperCase()); - } + if (fieldList == null || fieldList.isEmpty()) { + return blackFieldList == null || blackFieldList.isEmpty() + || !blackFieldList.contains(columnName.toUpperCase()); + } else { + return fieldList.contains(columnName.toUpperCase()); + } } public static TransactionBegin createTransactionBegin(long threadId) { @@ -1021,31 +1035,30 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar this.nameBlackFilter = nameBlackFilter; logger.warn("--> init table black filter : " + nameBlackFilter.toString()); } - + public void setFieldFilterMap(Map> fieldFilterMap) { - if (fieldFilterMap != null) { - this.fieldFilterMap = fieldFilterMap; - } else { - this.fieldFilterMap = new HashMap>(); - } - - - for (Map.Entry> entry : this.fieldFilterMap.entrySet()) { - logger.warn("--> init field filter : " + entry.getKey() + "->" + entry.getValue()); - } - } - + if (fieldFilterMap != null) { + this.fieldFilterMap = fieldFilterMap; + } else { + this.fieldFilterMap = new HashMap>(); + } + + for (Map.Entry> entry : this.fieldFilterMap.entrySet()) { + logger.warn("--> init field filter : " + entry.getKey() + "->" + entry.getValue()); + } + } + public void setFieldBlackFilterMap(Map> fieldBlackFilterMap) { - if (fieldBlackFilterMap != null) { - this.fieldBlackFilterMap = fieldBlackFilterMap; - } else { - this.fieldBlackFilterMap = new HashMap>(); - } - - for (Map.Entry> entry : this.fieldBlackFilterMap.entrySet()) { - logger.warn("--> init field black filter : " + entry.getKey() + "->" + entry.getValue()); - } - } + if (fieldBlackFilterMap != null) { + this.fieldBlackFilterMap = fieldBlackFilterMap; + } else { + this.fieldBlackFilterMap = new HashMap>(); + } + + for (Map.Entry> entry : this.fieldBlackFilterMap.entrySet()) { + logger.warn("--> init field black filter : " + entry.getKey() + "->" + entry.getValue()); + } + } public void setTableMetaCache(TableMetaCache tableMetaCache) { this.tableMetaCache = tableMetaCache; diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/RdsBinlogEventParserProxy.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/RdsBinlogEventParserProxy.java index 817d895f..53b770ef 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/RdsBinlogEventParserProxy.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/RdsBinlogEventParserProxy.java @@ -55,7 +55,7 @@ public class RdsBinlogEventParserProxy extends MysqlEventParser { rdsLocalBinlogEventParser.setLogPositionManager(this.getLogPositionManager()); rdsLocalBinlogEventParser.setDestination(destination); rdsLocalBinlogEventParser.setAlarmHandler(this.getAlarmHandler()); - rdsLocalBinlogEventParser.setConnectionCharset(this.connectionCharset); + rdsLocalBinlogEventParser.setConnectionCharsetStd(this.connectionCharset); rdsLocalBinlogEventParser.setConnectionCharsetNumber(this.connectionCharsetNumber); rdsLocalBinlogEventParser.setEnableTsdb(this.enableTsdb); rdsLocalBinlogEventParser.setEventBlackFilter(this.eventBlackFilter); diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/request/DescribeBackupPolicyRequest.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/request/DescribeBackupPolicyRequest.java index c5a9b501..cf51961a 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/request/DescribeBackupPolicyRequest.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/rds/request/DescribeBackupPolicyRequest.java @@ -31,7 +31,7 @@ public class DescribeBackupPolicyRequest extends AbstractRequest rollbackPosition.getTimestamp()) { continue; - } else if (rollbackPosition.getServerId().equals(snapshotPosition.getServerId()) + } else if (rollbackPosition.getServerId() == snapshotPosition.getServerId() && snapshotPosition.compareTo(rollbackPosition) > 0) { continue; } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta.java index 6f838b49..3fd4d112 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/MemoryTableMeta.java @@ -81,8 +81,7 @@ public class MemoryTableMeta implements TableMetaTSDB { && !StringUtils.startsWithIgnoreCase(StringUtils.trim(ddl), "alter user") && !StringUtils.startsWithIgnoreCase(StringUtils.trim(ddl), "drop user") && !StringUtils.startsWithIgnoreCase(StringUtils.trim(ddl), "create database")) { - String sql = trimSimpleComment(ddl); - repository.console(sql); + repository.console(ddl); } } catch (Throwable e) { logger.warn("parse faield : " + ddl, e); @@ -98,38 +97,6 @@ public class MemoryTableMeta implements TableMetaTSDB { return true; } - /** - * 过滤简单的行与块注释(--或# 开始的行 /*开始与*\/结束的块) - * - * @param sql - * @return - */ - private String trimSimpleComment(String sql) { - String[] sqls = sql.split("\n"); - if (sqls.length == 1) return sql; - - boolean isCommentBlock = false; - StringBuilder sb = new StringBuilder(sql.length()); - for (String psql : sqls) { - if (psql.startsWith("--") || psql.startsWith("#")) { - continue; - } else if (psql.trim().length() == 0) { - sb.append("\n"); - } - if (!isCommentBlock && !psql.trim().startsWith("/*") && !psql.trim().endsWith("*/")) { - sb.append(psql).append("\n"); - continue; - } - if (!isCommentBlock && psql.trim().startsWith("/*")) { - isCommentBlock = true; - } - if (isCommentBlock && psql.trim().endsWith("*/")) { - isCommentBlock = false; - } - } - return sb.toString(); - } - @Override public TableMeta find(String schema, String table) { List keys = Arrays.asList(schema, table); diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/DirectLogFetcherTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/DirectLogFetcherTest.java index 31caf4d8..a6d15892 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/DirectLogFetcherTest.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/DirectLogFetcherTest.java @@ -45,6 +45,7 @@ import com.taobao.tddl.dbsync.binlog.event.UpdateRowsLogEvent; import com.taobao.tddl.dbsync.binlog.event.WriteRowsLogEvent; import com.taobao.tddl.dbsync.binlog.event.XidLogEvent; import com.taobao.tddl.dbsync.binlog.event.mariadb.AnnotateRowsEvent; + @Ignore public class DirectLogFetcherTest { @@ -85,6 +86,9 @@ public class DirectLogFetcherTest { // event).getFilename(); System.out.println(((RotateLogEvent) event).getFilename()); break; + case LogEvent.TABLE_MAP_EVENT: + parseTableMapEvent((TableMapLogEvent) event); + break; case LogEvent.WRITE_ROWS_EVENT_V1: case LogEvent.WRITE_ROWS_EVENT: parseRowsEvent((WriteRowsLogEvent) event); @@ -173,13 +177,13 @@ public class DirectLogFetcherTest { logger.warn("update wait_timeout failed", e); } try { - update("set net_write_timeout=1800", connector); + update("set net_write_timeout=7200", connector); } catch (Exception e) { logger.warn("update net_write_timeout failed", e); } try { - update("set net_read_timeout=1800", connector); + update("set net_read_timeout=7200", connector); } catch (Exception e) { logger.warn("update net_read_timeout failed", e); } @@ -272,6 +276,18 @@ public class DirectLogFetcherTest { System.out.println("sql : " + new String(event.getRowsQuery().getBytes("ISO-8859-1"), charset.name())); } + public void parseTableMapEvent(TableMapLogEvent event) { + try { + String charsetDbName = new String(event.getDbName().getBytes("ISO-8859-1"), charset.name()); + event.setDbname(charsetDbName); + + String charsetTbName = new String(event.getTableName().getBytes("ISO-8859-1"), charset.name()); + event.setTblname(charsetTbName); + } catch (UnsupportedEncodingException e) { + throw new CanalParseException(e); + } + } + protected void parseXidEvent(XidLogEvent event) throws Exception { System.out.println(String.format("================> binlog[%s:%s]", binlogFileName, event.getHeader() .getLogPos() - event.getHeader().getEventLen())); diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/MysqlBinlogDumpPerformanceTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/MysqlBinlogDumpPerformanceTest.java index cfb5d132..e6214bd6 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/MysqlBinlogDumpPerformanceTest.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/MysqlBinlogDumpPerformanceTest.java @@ -1,10 +1,11 @@ package com.alibaba.otter.canal.parse; import java.net.InetSocketAddress; -import java.nio.charset.Charset; import java.util.List; import java.util.concurrent.atomic.AtomicLong; +import org.junit.Ignore; + import com.alibaba.otter.canal.common.AbstractCanalLifeCycle; import com.alibaba.otter.canal.parse.exception.CanalParseException; import com.alibaba.otter.canal.parse.inbound.mysql.MysqlEventParser; @@ -15,7 +16,6 @@ import com.alibaba.otter.canal.protocol.position.EntryPosition; import com.alibaba.otter.canal.protocol.position.LogPosition; import com.alibaba.otter.canal.sink.CanalEventSink; import com.alibaba.otter.canal.sink.exception.CanalSinkException; -import org.junit.Ignore; @Ignore public class MysqlBinlogDumpPerformanceTest { @@ -23,7 +23,7 @@ public class MysqlBinlogDumpPerformanceTest { public static void main(String args[]) { final MysqlEventParser controller = new MysqlEventParser(); final EntryPosition startPosition = new EntryPosition("mysql-bin.000007", 89796293L, 100L); - controller.setConnectionCharset(Charset.forName("UTF-8")); + controller.setConnectionCharset("UTF-8"); controller.setSlaveId(3344L); controller.setDetectingEnable(false); controller.setFilterQueryDml(true); diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/LocalBinlogDumpTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/LocalBinlogDumpTest.java index c21ec187..b38901f5 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/LocalBinlogDumpTest.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/LocalBinlogDumpTest.java @@ -31,7 +31,7 @@ public class LocalBinlogDumpTest { final EntryPosition startPosition = new EntryPosition("mysql-bin.000003", 123L); controller.setMasterInfo(new AuthenticationInfo(new InetSocketAddress("127.0.0.1", 3306), "canal", "canal")); - controller.setConnectionCharset(Charset.forName("UTF-8")); + controller.setConnectionCharsetStd(Charset.forName("UTF-8")); controller.setDirectory(directory); controller.setMasterPosition(startPosition); controller.setEventSink(new AbstractCanalEventSinkTest>() { diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlDumpTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlDumpTest.java index 212682b0..79bcfcf3 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlDumpTest.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlDumpTest.java @@ -32,7 +32,7 @@ public class MysqlDumpTest { final MysqlEventParser controller = new MysqlEventParser(); final EntryPosition startPosition = new EntryPosition("mysql-bin.000001", 4L); // startPosition.setGtid("f1ceb61a-a5d5-11e7-bdee-107c3dbcf8a7:1-17"); - controller.setConnectionCharset(Charset.forName("UTF-8")); + controller.setConnectionCharsetStd(Charset.forName("UTF-8")); controller.setSlaveId(3344L); controller.setDetectingEnable(false); controller.setMasterInfo(new AuthenticationInfo(new InetSocketAddress("127.0.0.1", 3306), "root", "hello")); diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/RdsLocalBinlogDumpTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/RdsLocalBinlogDumpTest.java index b718936f..943faf49 100644 --- a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/RdsLocalBinlogDumpTest.java +++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/RdsLocalBinlogDumpTest.java @@ -34,7 +34,7 @@ public class RdsLocalBinlogDumpTest { String directory = "/tmp/rds"; final RdsLocalBinlogEventParser controller = new RdsLocalBinlogEventParser(); controller.setMasterInfo(new AuthenticationInfo(new InetSocketAddress("127.0.0.1", 3306), "root", "hello")); - controller.setConnectionCharset(Charset.forName("UTF-8")); + controller.setConnectionCharsetStd(Charset.forName("UTF-8")); controller.setDirectory(directory); controller.setAccesskey(""); controller.setSecretkey(""); diff --git a/pom.xml b/pom.xml index 70f70b24..3ca37ebb 100644 --- a/pom.xml +++ b/pom.xml @@ -48,7 +48,7 @@ central - https://repo1.maven.org/maven2 + http://repo1.maven.org/maven2 true @@ -232,6 +232,11 @@ aviator 2.2.1 + + com.alibaba + fastjson + 1.2.58.sec06 + oro oro @@ -327,7 +332,7 @@ mysql mysql-connector-java - 5.1.47 + 5.1.48 diff --git a/prometheus/src/main/java/com/alibaba/otter/canal/prometheus/impl/PrometheusCanalEventDownStreamHandler.java b/prometheus/src/main/java/com/alibaba/otter/canal/prometheus/impl/PrometheusCanalEventDownStreamHandler.java index 0bfacffe..96e96944 100644 --- a/prometheus/src/main/java/com/alibaba/otter/canal/prometheus/impl/PrometheusCanalEventDownStreamHandler.java +++ b/prometheus/src/main/java/com/alibaba/otter/canal/prometheus/impl/PrometheusCanalEventDownStreamHandler.java @@ -1,13 +1,13 @@ package com.alibaba.otter.canal.prometheus.impl; +import java.util.List; +import java.util.concurrent.atomic.AtomicLong; + import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.CanalEntry.EntryType; import com.alibaba.otter.canal.sink.AbstractCanalEventDownStreamHandler; import com.alibaba.otter.canal.store.model.Event; -import java.util.List; -import java.util.concurrent.atomic.AtomicLong; - /** * @author Chuanyi Li */ @@ -26,17 +26,23 @@ public class PrometheusCanalEventDownStreamHandler extends AbstractCanalEventDow switch (type) { case TRANSACTIONBEGIN: { long exec = e.getExecuteTime(); - if (exec > 0) localExecTime = exec; + if (exec > 0) { + localExecTime = exec; + } break; } case ROWDATA: { long exec = e.getExecuteTime(); - if (exec > 0) localExecTime = exec; + if (exec > 0) { + localExecTime = exec; + } break; } case TRANSACTIONEND: { long exec = e.getExecuteTime(); - if (exec > 0) localExecTime = exec; + if (exec > 0) { + localExecTime = exec; + } transactionCounter.incrementAndGet(); break; } @@ -51,7 +57,7 @@ public class PrometheusCanalEventDownStreamHandler extends AbstractCanalEventDow } } if (localExecTime > 0) { - latestExecuteTime.lazySet(localExecTime); + latestExecuteTime.set(localExecTime); } } return events; diff --git a/server/src/main/java/com/alibaba/otter/canal/admin/CanalAdmin.java b/server/src/main/java/com/alibaba/otter/canal/admin/CanalAdmin.java index 69664ac9..9a5ed3fb 100644 --- a/server/src/main/java/com/alibaba/otter/canal/admin/CanalAdmin.java +++ b/server/src/main/java/com/alibaba/otter/canal/admin/CanalAdmin.java @@ -115,11 +115,4 @@ public interface CanalAdmin { * @return 日志信息 */ String instanceLog(String destination, String fileName, int lines); - - /** - * 获取meta - * @param destination - * @return - */ - String instanceMeta(String destination,String fileName); } diff --git a/server/src/main/java/com/alibaba/otter/canal/admin/handler/SessionHandler.java b/server/src/main/java/com/alibaba/otter/canal/admin/handler/SessionHandler.java index e4d7993b..ed5abea3 100644 --- a/server/src/main/java/com/alibaba/otter/canal/admin/handler/SessionHandler.java +++ b/server/src/main/java/com/alibaba/otter/canal/admin/handler/SessionHandler.java @@ -116,9 +116,6 @@ public class SessionHandler extends SimpleChannelHandler { message = canalAdmin.instanceLog(destination, file, count); } break; - case "meta": - message = canalAdmin.instanceMeta(destination, file); - break; default: byte[] errorBytes = AdminNettyUtils.errorPacket(301, MessageFormatter.format("LogAdmin type={} is unknown", type).getMessage()); diff --git a/server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java b/server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java index dae91df2..50be0d45 100644 --- a/server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java +++ b/server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java @@ -7,22 +7,21 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import com.alibaba.otter.canal.connector.core.util.Callback; -import com.alibaba.otter.canal.connector.core.producer.MQDestination; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.slf4j.MDC; +import com.alibaba.otter.canal.connector.core.config.MQProperties; +import com.alibaba.otter.canal.connector.core.producer.MQDestination; import com.alibaba.otter.canal.connector.core.spi.CanalMQProducer; +import com.alibaba.otter.canal.connector.core.util.Callback; import com.alibaba.otter.canal.instance.core.CanalInstance; import com.alibaba.otter.canal.instance.core.CanalMQConfig; import com.alibaba.otter.canal.protocol.ClientIdentity; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.server.embedded.CanalServerWithEmbedded; -import com.alibaba.otter.canal.connector.core.config.MQProperties; - public class CanalMQStarter { private static final Logger logger = LoggerFactory.getLogger(CanalMQStarter.class); @@ -170,8 +169,10 @@ public class CanalMQStarter { while (running && destinationRunning.get()) { Message message; if (getTimeout != null && getTimeout > 0) { - message = canalServer - .getWithoutAck(clientIdentity, getBatchSize, getTimeout.longValue(), TimeUnit.MILLISECONDS); + message = canalServer.getWithoutAck(clientIdentity, + getBatchSize, + getTimeout.longValue(), + TimeUnit.MILLISECONDS); } else { message = canalServer.getWithoutAck(clientIdentity, getBatchSize); } diff --git a/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java b/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java index 27df9074..4c719ec1 100644 --- a/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java +++ b/sink/src/main/java/com/alibaba/otter/canal/sink/entry/EntryEventSink.java @@ -104,8 +104,12 @@ public class EntryEventSink extends AbstractCanalEventSink