diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java index 3574b678..00939c58 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/CanalClientConfig.java @@ -33,7 +33,14 @@ public class CanalClientConfig { // aliyun ak/sk private String accessKey; private String secretKey; - + // 是否启用消息轨迹 + private boolean enableMessageTrace; + // 在使用阿里云商业化mq服务时,如果想使用云上消息轨迹功能,请设置此配置为true + private String accessChannel; + // 用于使用开源RocketMQ时,设置自定义的消息轨迹topic + private String customizedTraceTopic; + // 开源RocketMQ命名空间 + private String namespace; // canal adapters 配置 private List canalAdapters; @@ -133,6 +140,38 @@ public class CanalClientConfig { this.canalAdapters = canalAdapters; } + public boolean isEnableMessageTrace() { + return enableMessageTrace; + } + + public void setEnableMessageTrace(boolean enableMessageTrace) { + this.enableMessageTrace = enableMessageTrace; + } + + public String getAccessChannel() { + return accessChannel; + } + + public void setAccessChannel(String accessChannel) { + this.accessChannel = accessChannel; + } + + public String getCustomizedTraceTopic() { + return customizedTraceTopic; + } + + public void setCustomizedTraceTopic(String customizedTraceTopic) { + this.customizedTraceTopic = customizedTraceTopic; + } + + public String getNamespace() { + return namespace; + } + + public void setNamespace(String namespace) { + this.namespace = namespace; + } + public static class CanalAdapter { private String instance; // 实例名 diff --git a/client-adapter/launcher/pom.xml b/client-adapter/launcher/pom.xml index 1aa095b5..e93956c9 100644 --- a/client-adapter/launcher/pom.xml +++ b/client-adapter/launcher/pom.xml @@ -62,7 +62,7 @@ org.apache.rocketmq rocketmq-client - 4.3.0 + 4.5.1 diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterLoader.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterLoader.java index 9974928d..7623c048 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterLoader.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterLoader.java @@ -130,7 +130,11 @@ public class CanalAdapterLoader { canalOuterAdapterGroups, canalClientConfig.getAccessKey(), canalClientConfig.getSecretKey(), - canalClientConfig.getFlatMessage()); + canalClientConfig.getFlatMessage(), + canalClientConfig.isEnableMessageTrace(), + canalClientConfig.getCustomizedTraceTopic(), + canalClientConfig.getAccessChannel(), + canalClientConfig.getNamespace()); canalMQWorker.put(canalAdapter.getInstance() + "-rocketmq-" + group.getGroupId(), rocketMQWorker); rocketMQWorker.start(); diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java index 5da53763..3193553c 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/loader/CanalAdapterRocketMQWorker.java @@ -40,6 +40,30 @@ public class CanalAdapterRocketMQWorker extends AbstractCanalAdapterWorker { logger.info("RocketMQ consumer config topic:{}, nameServer:{}, groupId:{}", topic, nameServers, groupId); } + public CanalAdapterRocketMQWorker(CanalClientConfig canalClientConfig, String nameServers, String topic, + String groupId, List> canalOuterAdapters, String accessKey, + String secretKey, boolean flatMessage, boolean enableMessageTrace, + String customizedTraceTopic, String accessChannel, String namespace) { + super(canalOuterAdapters); + this.canalClientConfig = canalClientConfig; + this.topic = topic; + this.flatMessage = flatMessage; + super.canalDestination = topic; + super.groupId = groupId; + this.connector = new RocketMQCanalConnector(nameServers, + topic, + groupId, + accessKey, + secretKey, + canalClientConfig.getBatchSize(), + flatMessage, + enableMessageTrace, + customizedTraceTopic, + accessChannel, + namespace); + logger.info("RocketMQ consumer config topic:{}, nameServer:{}, groupId:{}", topic, nameServers, groupId); + } + @Override protected void process() { while (!running) { diff --git a/client-adapter/launcher/src/main/resources/application.yml b/client-adapter/launcher/src/main/resources/application.yml index b53e8735..14bc3e6e 100644 --- a/client-adapter/launcher/src/main/resources/application.yml +++ b/client-adapter/launcher/src/main/resources/application.yml @@ -10,7 +10,7 @@ canal.conf: mode: tcp # kafka rocketMQ canalServerHost: 127.0.0.1:11111 # zookeeperHosts: slave1:2181 -# mqServers: 127.0.0.1:9092 #or rocketmq +# mqServers: 127.0.0.1:9092 #or rocketmq nameservers # flatMessage: true batchSize: 500 syncBatchSize: 1000 @@ -18,6 +18,11 @@ canal.conf: timeout: accessKey: secretKey: +# enableMessageTrace: +# accessChannel: +# customizedTraceTopic: +# namespace: + # srcDataSources: # defaultDS: # url: jdbc:mysql://127.0.0.1:3306/mytest?useUnicode=true 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 a34bb784..350b29e9 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 @@ -272,7 +272,7 @@ public class RdbSyncService { batchExecutor.execute(insertSql.toString(), values); } catch (SQLException e) { if (skipDupException - && (e.getMessage().contains("Duplicate entry") || e.getMessage().startsWith("ORA-00001: 违反唯一约束条件"))) { + && (e.getMessage().contains("Duplicate entry") || e.getMessage().startsWith("ORA-00001:"))) { // ignore // TODO 增加更多关系数据库的主键冲突的错误码 } else { diff --git a/client/pom.xml b/client/pom.xml index eeb62ea0..13083413 100644 --- a/client/pom.xml +++ b/client/pom.xml @@ -106,8 +106,8 @@ rocketmq-client - com.aliyun.openservices - aliware-apache-rocketmq-cloud + org.apache.rocketmq + rocketmq-acl diff --git a/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java b/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java index a39c1ba5..a1a6d83b 100644 --- a/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java +++ b/client/src/main/java/com/alibaba/otter/canal/client/rocketmq/RocketMQCanalConnector.java @@ -6,6 +6,9 @@ import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import org.apache.commons.lang.StringUtils; +import org.apache.rocketmq.acl.common.AclClientRPCHook; +import org.apache.rocketmq.acl.common.SessionCredentials; +import org.apache.rocketmq.client.AccessChannel; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus; @@ -24,8 +27,6 @@ import com.alibaba.otter.canal.client.impl.SimpleCanalConnector; import com.alibaba.otter.canal.protocol.FlatMessage; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.exception.CanalClientException; -import com.aliyun.openservices.apache.api.impl.authority.SessionCredentials; -import com.aliyun.openservices.apache.api.impl.rocketmq.ClientRPCHook; import com.google.common.collect.Lists; /** @@ -40,7 +41,8 @@ import com.google.common.collect.Lists; */ public class RocketMQCanalConnector implements CanalMQConnector { - private static final Logger logger = LoggerFactory.getLogger(RocketMQCanalConnector.class); + private static final Logger logger = LoggerFactory.getLogger(RocketMQCanalConnector.class); + private static final String CLOUD_ACCESS_CHANNEL = "cloud"; private String nameServer; private String topic; @@ -54,6 +56,26 @@ public class RocketMQCanalConnector implements CanalMQConnector { private volatile ConsumerBatchMessage lastGetBatchMessage = null; private String accessKey; private String secretKey; + private String customizedTraceTopic; + private boolean enableMessageTrace = false; + private String accessChannel; + private String namespace; + + public RocketMQCanalConnector(String nameServer, String topic, String groupName, String accessKey, + String secretKey, Integer batchSize, boolean flatMessage, boolean enableMessageTrace, + String customizedTraceTopic, String accessChannel, String namespace) { + this(nameServer, topic, groupName, accessKey, secretKey, batchSize, flatMessage, enableMessageTrace, customizedTraceTopic, accessChannel); + this.namespace = namespace; + } + + public RocketMQCanalConnector(String nameServer, String topic, String groupName, String accessKey, + String secretKey, Integer batchSize, boolean flatMessage, boolean enableMessageTrace, + String customizedTraceTopic, String accessChannel) { + this(nameServer, topic, groupName, accessKey, secretKey, batchSize, flatMessage); + this.enableMessageTrace = enableMessageTrace; + this.customizedTraceTopic = customizedTraceTopic; + this.accessChannel = accessChannel; + } public RocketMQCanalConnector(String nameServer, String topic, String groupName, Integer batchSize, boolean flatMessage){ @@ -73,16 +95,24 @@ public class RocketMQCanalConnector implements CanalMQConnector { } public void connect() throws CanalClientException { - RPCHook rpcHook = null; if (null != accessKey && accessKey.length() > 0 && null != secretKey && secretKey.length() > 0) { SessionCredentials sessionCredentials = new SessionCredentials(); sessionCredentials.setAccessKey(accessKey); sessionCredentials.setSecretKey(secretKey); - rpcHook = new ClientRPCHook(sessionCredentials); + rpcHook = new AclClientRPCHook(sessionCredentials); } - rocketMQConsumer = new DefaultMQPushConsumer(groupName, rpcHook, new AllocateMessageQueueAveragely()); + + rocketMQConsumer = new DefaultMQPushConsumer(groupName, rpcHook, new AllocateMessageQueueAveragely(), enableMessageTrace, customizedTraceTopic); rocketMQConsumer.setVipChannelEnabled(false); + if (CLOUD_ACCESS_CHANNEL.equals(this.accessChannel)) { + rocketMQConsumer.setAccessChannel(AccessChannel.CLOUD); + } + + if (!StringUtils.isEmpty(this.namespace)) { + rocketMQConsumer.setNamespace(this.namespace); + } + if (!StringUtils.isBlank(nameServer)) { rocketMQConsumer.setNamesrvAddr(nameServer); } @@ -131,7 +161,9 @@ public class RocketMQCanalConnector implements CanalMQConnector { } private boolean process(List messageExts) { - logger.info("Get Message:{}", messageExts); + if (logger.isDebugEnabled()) { + logger.debug("Get Message: {}", messageExts); + } List messageList = Lists.newArrayList(); for (MessageExt messageExt : messageExts) { byte[] data = messageExt.getBody(); diff --git a/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/QueryLogEvent.java b/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/QueryLogEvent.java index c732da99..8cc49795 100644 --- a/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/QueryLogEvent.java +++ b/dbsync/src/main/java/com/taobao/tddl/dbsync/binlog/event/QueryLogEvent.java @@ -706,10 +706,8 @@ public class QueryLogEvent extends LogEvent { * That's why you must write status vars in growing * order of code */ - if (logger.isDebugEnabled()) { - logger.debug("Query_log_event has unknown status vars (first has code: " + code - + "), skipping the rest of them"); - } + logger.error("Query_log_event has unknown status vars (first has code: " + code + + "), skipping the rest of them"); break; // Break loop } } diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java index 45cffa5a..6d0d9794 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalConstants.java @@ -10,49 +10,56 @@ import java.text.MessageFormat; */ public class CanalConstants { - public static final String MDC_DESTINATION = "destination"; - public static final String ROOT = "canal"; - public static final String CANAL_ID = ROOT + "." + "id"; - public static final String CANAL_IP = ROOT + "." + "ip"; - public static final String CANAL_PORT = ROOT + "." + "port"; - public static final String CANAL_METRICS_PULL_PORT = ROOT + "." + "metrics.pull.port"; - public static final String CANAL_ZKSERVERS = ROOT + "." + "zkServers"; - public static final String CANAL_WITHOUT_NETTY = ROOT + "." + "withoutNetty"; + public static final String MDC_DESTINATION = "destination"; + public static final String ROOT = "canal"; + public static final String CANAL_ID = ROOT + "." + "id"; + public static final String CANAL_IP = ROOT + "." + "ip"; + public static final String CANAL_PORT = ROOT + "." + "port"; + public static final String CANAL_METRICS_PULL_PORT = ROOT + "." + "metrics.pull.port"; + public static final String CANAL_ZKSERVERS = ROOT + "." + "zkServers"; + public static final String CANAL_WITHOUT_NETTY = ROOT + "." + "withoutNetty"; - public static final String CANAL_DESTINATIONS = ROOT + "." + "destinations"; - public static final String CANAL_AUTO_SCAN = ROOT + "." + "auto.scan"; - public static final String CANAL_AUTO_SCAN_INTERVAL = ROOT + "." + "auto.scan.interval"; - public static final String CANAL_CONF_DIR = ROOT + "." + "conf.dir"; - public static final String CANAL_SERVER_MODE = ROOT + "." + "serverMode"; + public static final String CANAL_DESTINATIONS = ROOT + "." + "destinations"; + public static final String CANAL_AUTO_SCAN = ROOT + "." + "auto.scan"; + public static final String CANAL_AUTO_SCAN_INTERVAL = ROOT + "." + "auto.scan.interval"; + public static final String CANAL_CONF_DIR = ROOT + "." + "conf.dir"; + public static final String CANAL_SERVER_MODE = ROOT + "." + "serverMode"; - public static final String CANAL_DESTINATION_SPLIT = ","; - public static final String GLOBAL_NAME = "global"; + public static final String CANAL_DESTINATION_SPLIT = ","; + public static final String GLOBAL_NAME = "global"; - public static final String INSTANCE_MODE_TEMPLATE = ROOT + "." + "instance.{0}.mode"; - public static final String INSTANCE_LAZY_TEMPLATE = ROOT + "." + "instance.{0}.lazy"; - public static final String INSTANCE_MANAGER_ADDRESS_TEMPLATE = ROOT + "." + "instance.{0}.manager.address"; - public static final String INSTANCE_SPRING_XML_TEMPLATE = ROOT + "." + "instance.{0}.spring.xml"; + public static final String INSTANCE_MODE_TEMPLATE = ROOT + "." + "instance.{0}.mode"; + public static final String INSTANCE_LAZY_TEMPLATE = ROOT + "." + "instance.{0}.lazy"; + public static final String INSTANCE_MANAGER_ADDRESS_TEMPLATE = ROOT + "." + "instance.{0}.manager.address"; + public static final String INSTANCE_SPRING_XML_TEMPLATE = ROOT + "." + "instance.{0}.spring.xml"; - public static final String CANAL_DESTINATION_PROPERTY = ROOT + ".instance.destination"; + public static final String CANAL_DESTINATION_PROPERTY = ROOT + ".instance.destination"; - public static final String CANAL_SOCKETCHANNEL = ROOT + "." + "socketChannel"; + public static final String CANAL_SOCKETCHANNEL = ROOT + "." + "socketChannel"; - public static final String CANAL_MQ_SERVERS = ROOT + "." + "mq.servers"; - public static final String CANAL_MQ_RETRIES = ROOT + "." + "mq.retries"; - public static final String CANAL_MQ_BATCHSIZE = ROOT + "." + "mq.batchSize"; - public static final String CANAL_MQ_LINGERMS = ROOT + "." + "mq.lingerMs"; - public static final String CANAL_MQ_MAXREQUESTSIZE = ROOT + "." + "mq.maxRequestSize"; - public static final String CANAL_MQ_BUFFERMEMORY = ROOT + "." + "mq.bufferMemory"; - public static final String CANAL_MQ_CANALBATCHSIZE = ROOT + "." + "mq.canalBatchSize"; - public static final String CANAL_MQ_CANALGETTIMEOUT = ROOT + "." + "mq.canalGetTimeout"; - public static final String CANAL_MQ_FLATMESSAGE = ROOT + "." + "mq.flatMessage"; - public static final String CANAL_MQ_COMPRESSION_TYPE = ROOT + "." + "mq.compressionType"; - public static final String CANAL_MQ_ACKS = ROOT + "." + "mq.acks"; - public static final String CANAL_MQ_TRANSACTION = ROOT + "." + "mq.transaction"; - public static final String CANAL_MQ_PRODUCERGROUP = ROOT + "." + "mq.producerGroup"; - public static final String CANAL_ALIYUN_ACCESSKEY = ROOT + "." + "aliyun.accessKey"; - public static final String CANAL_ALIYUN_SECRETKEY = ROOT + "." + "aliyun.secretKey"; - public static final String CANAL_MQ_PROPERTIES = ROOT + "." + "mq.properties"; + public static final String CANAL_MQ_SERVERS = ROOT + "." + "mq.servers"; + public static final String CANAL_MQ_RETRIES = ROOT + "." + "mq.retries"; + public static final String CANAL_MQ_BATCHSIZE = ROOT + "." + "mq.batchSize"; + public static final String CANAL_MQ_LINGERMS = ROOT + "." + "mq.lingerMs"; + public static final String CANAL_MQ_MAXREQUESTSIZE = ROOT + "." + "mq.maxRequestSize"; + public static final String CANAL_MQ_BUFFERMEMORY = ROOT + "." + "mq.bufferMemory"; + public static final String CANAL_MQ_CANALBATCHSIZE = ROOT + "." + "mq.canalBatchSize"; + public static final String CANAL_MQ_CANALGETTIMEOUT = ROOT + "." + "mq.canalGetTimeout"; + public static final String CANAL_MQ_FLATMESSAGE = ROOT + "." + "mq.flatMessage"; + public static final String CANAL_MQ_COMPRESSION_TYPE = ROOT + "." + "mq.compressionType"; + public static final String CANAL_MQ_ACKS = ROOT + "." + "mq.acks"; + public static final String CANAL_MQ_TRANSACTION = ROOT + "." + "mq.transaction"; + public static final String CANAL_MQ_PRODUCERGROUP = ROOT + "." + "mq.producerGroup"; + public static final String CANAL_ALIYUN_ACCESSKEY = ROOT + "." + "aliyun.accessKey"; + public static final String CANAL_ALIYUN_SECRETKEY = ROOT + "." + "aliyun.secretKey"; + public static final String CANAL_MQ_PROPERTIES = ROOT + "." + "mq.properties"; + public static final String CANAL_MQ_ENABLE_MESSAGE_TRACE = ROOT + "." + "mq.enableMessageTrace"; + public static final String CANAL_MQ_ACCESS_CHANNEL = ROOT + "." + "mq.accessChannel"; + public static final String CANAL_MQ_CUSTOMIZED_TRACE_TOPIC = ROOT + "." + "mq.customizedTraceTopic"; + public static final String CANAL_MQ_NAMESPACE = ROOT + "." + "mq.namespace"; + public static final String CANAL_MQ_KAFKA_KERBEROS_ENABLE = ROOT + "." + "mq.kafka.kerberos.enable"; + public static final String CANAL_MQ_KAFKA_KERBEROS_KRB5FILEPATH = ROOT + "." + "mq.kafka.kerberos.krb5FilePath"; + public static final String CANAL_MQ_KAFKA_KERBEROS_JAASFILEPATH = ROOT + "." + "mq.kafka.kerberos.jaasFilePath"; public static String getInstanceModeKey(String destination) { return MessageFormat.format(INSTANCE_MODE_TEMPLATE, destination); 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 d8ca17fd..c0692227 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 @@ -152,6 +152,9 @@ public class CanalController { try { MDC.put(CanalConstants.MDC_DESTINATION, String.valueOf(destination)); embededCanalServer.start(destination); + if (canalMQStarter != null) { + canalMQStarter.startDestination(destination); + } } finally { MDC.remove(CanalConstants.MDC_DESTINATION); } @@ -160,6 +163,9 @@ public class CanalController { public void processActiveExit() { try { MDC.put(CanalConstants.MDC_DESTINATION, String.valueOf(destination)); + if (canalMQStarter != null) { + canalMQStarter.stopDestination(destination); + } embededCanalServer.stop(destination); } finally { MDC.remove(CanalConstants.MDC_DESTINATION); @@ -234,9 +240,6 @@ public class CanalController { ServerRunningMonitor runningMonitor = ServerRunningMonitors.getRunningMonitor(destination); if (!config.getLazy() && !runningMonitor.isStart()) { runningMonitor.start(); - if (canalMQStarter != null) { - canalMQStarter.startDestination(destination); - } } } } @@ -245,9 +248,6 @@ public class CanalController { // 此处的stop,代表强制退出,非HA机制,所以需要退出HA的monitor和配置信息 InstanceConfig config = instanceConfigs.remove(destination); if (config != null) { - if (canalMQStarter != null) { - canalMQStarter.stopDestination(destination); - } embededCanalServer.stop(destination); ServerRunningMonitor runningMonitor = ServerRunningMonitors.getRunningMonitor(destination); if (runningMonitor.isStart()) { diff --git a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalStater.java b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalStater.java index 3a465bb0..09a227be 100644 --- a/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalStater.java +++ b/deployer/src/main/java/com/alibaba/otter/canal/deployer/CanalStater.java @@ -1,15 +1,5 @@ package com.alibaba.otter.canal.deployer; -import java.io.File; -import java.io.FileFilter; -import java.util.Arrays; -import java.util.List; -import java.util.Properties; - -import org.apache.commons.lang.StringUtils; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import com.alibaba.otter.canal.common.MQProperties; import com.alibaba.otter.canal.kafka.CanalKafkaProducer; import com.alibaba.otter.canal.rocketmq.CanalRocketMQProducer; @@ -18,6 +8,15 @@ import com.alibaba.otter.canal.spi.CanalMQProducer; import com.google.common.base.Function; import com.google.common.base.Joiner; import com.google.common.collect.Lists; +import org.apache.commons.lang.StringUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.File; +import java.io.FileFilter; +import java.util.Arrays; +import java.util.List; +import java.util.Properties; /** * Canal server 启动类 @@ -207,6 +206,41 @@ public class CanalStater { mqProperties.setProducerGroup(producerGroup); } + String enableMessageTrace = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_ENABLE_MESSAGE_TRACE); + if (!StringUtils.isEmpty(enableMessageTrace)) { + mqProperties.setEnableMessageTrace(Boolean.valueOf(enableMessageTrace)); + } + + String accessChannel = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_ACCESS_CHANNEL); + if (!StringUtils.isEmpty(accessChannel)) { + mqProperties.setAccessChannel(accessChannel); + } + + String customizedTraceTopic = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_CUSTOMIZED_TRACE_TOPIC); + if (!StringUtils.isEmpty(customizedTraceTopic)) { + mqProperties.setCustomizedTraceTopic(customizedTraceTopic); + } + + String namespace = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_NAMESPACE); + if (!StringUtils.isEmpty(namespace)) { + mqProperties.setNamespace(namespace); + } + + String kafkaKerberosEnable = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_KAFKA_KERBEROS_ENABLE); + if (!StringUtils.isEmpty(kafkaKerberosEnable)) { + mqProperties.setKerberosEnable(Boolean.valueOf(kafkaKerberosEnable)); + } + + String kafkaKerberosKrb5Filepath = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_KAFKA_KERBEROS_KRB5FILEPATH); + if (!StringUtils.isEmpty(kafkaKerberosKrb5Filepath)) { + mqProperties.setKerberosKrb5FilePath(kafkaKerberosKrb5Filepath); + } + + String kafkaKerberosJaasFilepath = CanalController.getProperty(properties, CanalConstants.CANAL_MQ_KAFKA_KERBEROS_JAASFILEPATH); + if (!StringUtils.isEmpty(kafkaKerberosJaasFilepath)) { + mqProperties.setKerberosJaasFilePath(kafkaKerberosJaasFilepath); + } + for (Object key : properties.keySet()) { key = StringUtils.trim(key.toString()); if (((String) key).startsWith(CanalConstants.CANAL_MQ_PROPERTIES)) { diff --git a/deployer/src/main/resources/canal.properties b/deployer/src/main/resources/canal.properties index 3ba5404c..106005f1 100644 --- a/deployer/src/main/resources/canal.properties +++ b/deployer/src/main/resources/canal.properties @@ -118,3 +118,15 @@ canal.mq.acks = all # use transaction for kafka flatMessage batch produce canal.mq.transaction = true #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 = + +################################################## +######### Kafka Kerberos Info ############# +################################################## +canal.mq.kafka.kerberos.enable = false +canal.mq.kafka.kerberos.krb5FilePath = "../conf/kerberos/krb5.conf" +canal.mq.kafka.kerberos.jaasFilePath = "../conf/kerberos/jaas.conf" \ No newline at end of file diff --git a/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml b/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml index c05b5200..6c34241d 100644 --- a/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml +++ b/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml @@ -43,6 +43,7 @@ + diff --git a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java index 55672db1..f9952e61 100644 --- a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java +++ b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java @@ -7,5 +7,9 @@ public abstract class AbstractRocektMQTest extends BaseCanalClientTest { public static String topic = "example"; public static String groupId = "group"; public static String nameServers = "localhost:9876"; - + public static String accessKey = ""; + public static String secretKey = ""; + public static boolean enableMessageTrace = false; + public static String accessChannel = "local"; + public static String namespace = ""; } diff --git a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java index 2a6c5c17..3c88900e 100644 --- a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java +++ b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java @@ -1,16 +1,13 @@ package com.alibaba.otter.canal.example.rocketmq; +import com.alibaba.otter.canal.client.rocketmq.RocketMQCanalConnector; +import com.alibaba.otter.canal.protocol.Message; import java.util.List; import java.util.concurrent.TimeUnit; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.util.Assert; -import com.alibaba.otter.canal.client.rocketmq.RocketMQCanalConnector; -import com.alibaba.otter.canal.example.kafka.AbstractKafkaTest; -import com.alibaba.otter.canal.protocol.Message; - /** * Kafka client example * @@ -34,15 +31,27 @@ public class CanalRocketMQClientExample extends AbstractRocektMQTest { } }; - public CanalRocketMQClientExample(String nameServers, String topic, String groupId){ + public CanalRocketMQClientExample(String nameServers, String topic, String groupId) { connector = new RocketMQCanalConnector(nameServers, topic, groupId, 500, false); } + public CanalRocketMQClientExample(String nameServers, String topic, String groupId, boolean enableMessageTrace, + String accessKey, String secretKey, String accessChannel, String namespace) { + connector = new RocketMQCanalConnector(nameServers, topic, groupId, accessKey, + secretKey, -1, false, enableMessageTrace, + null, accessChannel, namespace); + } + public static void main(String[] args) { try { final CanalRocketMQClientExample rocketMQClientExample = new CanalRocketMQClientExample(nameServers, topic, - groupId); + groupId, + enableMessageTrace, + accessKey, + secretKey, + accessChannel, + namespace); logger.info("## Start the rocketmq consumer: {}-{}", topic, groupId); rocketMQClientExample.start(); logger.info("## The canal rocketmq consumer is running now ......"); @@ -108,7 +117,7 @@ public class CanalRocketMQClientExample extends AbstractRocektMQTest { connector.connect(); connector.subscribe(); while (running) { - List messages = connector.getListWithoutAck(100L, TimeUnit.MILLISECONDS); // 获取message + List messages = connector.getListWithoutAck(1000L, TimeUnit.MILLISECONDS); // 获取message for (Message message : messages) { long batchId = message.getId(); int size = message.getEntries().size(); diff --git a/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java b/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java index b1b64d5b..6993a55a 100644 --- a/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java +++ b/instance/core/src/main/java/com/alibaba/otter/canal/instance/core/AbstractCanalInstance.java @@ -53,10 +53,14 @@ public class AbstractCanalInstance extends AbstractCanalLifeCycle implements Can // 处理group的模式 List eventParsers = ((GroupEventParser) eventParser).getEventParsers(); for (CanalEventParser singleEventParser : eventParsers) {// 需要遍历启动 - ((AbstractEventParser) singleEventParser).setEventFilter(aviaterFilter); + if(singleEventParser instanceof AbstractEventParser) { + ((AbstractEventParser) singleEventParser).setEventFilter(aviaterFilter); + } } } else { - ((AbstractEventParser) eventParser).setEventFilter(aviaterFilter); + if(eventParser instanceof AbstractEventParser) { + ((AbstractEventParser) eventParser).setEventFilter(aviaterFilter); + } } } 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 43eaaf33..49030fc1 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 @@ -1,6 +1,7 @@ package com.alibaba.otter.canal.instance.manager; import java.io.File; +import java.io.FilenameFilter; import java.net.InetSocketAddress; import java.net.URL; import java.net.URLClassLoader; @@ -117,7 +118,13 @@ public class CanalInstanceWithManager extends AbstractCanalInstance { } else { try { File externalLibDir = new File(alarmHandlerPluginDir); - File[] jarFiles = externalLibDir.listFiles((dir1, name) -> name.endsWith(".jar")); + File[] jarFiles = externalLibDir.listFiles(new FilenameFilter() { + + @Override + public boolean accept(File dir, String name) { + return name.endsWith(".jar"); + } + }); if (jarFiles == null || jarFiles.length == 0) { throw new IllegalStateException(String.format("alarmHandlerPluginDir [%s] can't find any name endswith \".jar\" file.", alarmHandlerPluginDir)); @@ -126,14 +133,16 @@ public class CanalInstanceWithManager extends AbstractCanalInstance { for (int i = 0; i < jarFiles.length; i++) { urls[i] = jarFiles[i].toURI().toURL(); } - ClassLoader currentClassLoader = new URLClassLoader(urls, CanalInstanceWithManager.class.getClassLoader()); - Class _alarmClass = - (Class)currentClassLoader.loadClass(alarmHandlerClass); + ClassLoader currentClassLoader = new URLClassLoader(urls, + CanalInstanceWithManager.class.getClassLoader()); + Class _alarmClass = (Class) currentClassLoader.loadClass(alarmHandlerClass); alarmHandler = _alarmClass.newInstance(); logger.info("init [{}] alarm handler success.", alarmHandlerClass); } catch (Throwable e) { String errorMsg = String.format("init alarmHandlerPluginDir [%s] alarm handler [%s] error: %s", - alarmHandlerPluginDir, alarmHandlerClass, ExceptionUtils.getFullStackTrace(e)); + alarmHandlerPluginDir, + alarmHandlerClass, + ExceptionUtils.getFullStackTrace(e)); logger.error(errorMsg); throw new CanalException(errorMsg, e); } diff --git a/logo.png b/logo.png new file mode 100644 index 00000000..3e2cc642 Binary files /dev/null and b/logo.png differ diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/AbstractEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/AbstractEventParser.java index 73b35df6..409c15d2 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/AbstractEventParser.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/AbstractEventParser.java @@ -242,7 +242,7 @@ public abstract class AbstractEventParser extends AbstractCanalLifeCycle if (parallel) { // build stage processor multiStageCoprocessor = buildMultiStageCoprocessor(); - if (isGTIDMode()) { + if (isGTIDMode() && StringUtils.isNotEmpty(startPosition.getGtid())) { // 判断所属instance是否启用GTID模式,是的话调用ErosaConnection中GTID对应方法dump数据 GTIDSet gtidSet = MysqlGTIDSet.parse(startPosition.getGtid()); ((MysqlMultiStageCoprocessor) multiStageCoprocessor).setGtidSet(gtidSet); @@ -260,7 +260,7 @@ public abstract class AbstractEventParser extends AbstractCanalLifeCycle } } } else { - if (isGTIDMode()) { + if (isGTIDMode() && StringUtils.isNotEmpty(startPosition.getGtid())) { // 判断所属instance是否启用GTID模式,是的话调用ErosaConnection中GTID对应方法dump数据 erosaConnection.dump(MysqlGTIDSet.parse(startPosition.getGtid()), sinkHandler); } else { diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java index d2ce8356..68eb5062 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java @@ -355,11 +355,14 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE // GTID模式下,CanalLogPositionManager里取最后的gtid,没有则取instanc配置中的 LogPosition logPosition = getLogPositionManager().getLatestIndexBy(destination); if (logPosition != null) { - return logPosition.getPostion(); - } - - if (masterPosition != null && StringUtils.isNotEmpty(masterPosition.getGtid())) { - return masterPosition; + // 如果以前是非GTID模式,后来调整为了GTID模式,那么为了保持兼容,需要判断gtid是否为空 + if (StringUtils.isNotEmpty(logPosition.getPostion().getGtid())) { + return logPosition.getPostion(); + } + }else { + if (masterPosition != null && StringUtils.isNotEmpty(masterPosition.getGtid())) { + return masterPosition; + } } } diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlMultiStageCoprocessor.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlMultiStageCoprocessor.java index 70aba4f1..5c5ad87a 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlMultiStageCoprocessor.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlMultiStageCoprocessor.java @@ -441,6 +441,8 @@ public class MysqlMultiStageCoprocessor extends AbstractCanalLifeCycle implement @Override public void handleEventException(final Throwable ex, final long sequence, final Object event) { + //异常上抛,否则processEvents的逻辑会默认会mark为成功执行,有丢数据风险 + throw new CanalParseException(ex); } @Override diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java index 25554c29..8c9d0e39 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/TableMetaCache.java @@ -81,7 +81,7 @@ public class TableMetaCache { } } - private TableMeta getTableMetaByDB(String fullname) throws IOException { + private synchronized TableMeta getTableMetaByDB(String fullname) throws IOException { try { ResultSetPacket packet = connection.query("show create table " + fullname); String[] names = StringUtils.split(fullname, "`.`"); @@ -159,7 +159,7 @@ public class TableMetaCache { return getTableMeta(schema, table, true, position); } - public TableMeta getTableMeta(String schema, String table, boolean useCache, EntryPosition position) { + public synchronized TableMeta getTableMeta(String schema, String table, boolean useCache, EntryPosition position) { TableMeta tableMeta = null; if (tableMetaTSDB != null) { tableMeta = tableMetaTSDB.find(schema, table); diff --git a/pom.xml b/pom.xml index 28a0ca4c..5f875209 100644 --- a/pom.xml +++ b/pom.xml @@ -99,6 +99,7 @@ 1.8 UTF-8 3.2.18.RELEASE + 4.5.1 0.8.3 2.22.1 -server -Xms512m -Xmx1024m -Dfile.encoding=UTF-8 @@ -252,7 +253,7 @@ com.alibaba.fastsql fastsql - 2.0.0_preview_855 + 2.0.0_preview_914 com.alibaba @@ -309,14 +310,14 @@ 3.0.2 - com.aliyun.openservices - aliware-apache-rocketmq-cloud - 1.0 + org.apache.rocketmq + rocketmq-client + ${rocketmq_version} org.apache.rocketmq - rocketmq-client - 4.3.0 + rocketmq-acl + ${rocketmq_version} javax.annotation diff --git a/server/pom.xml b/server/pom.xml index 4a7eb6ef..c1061312 100644 --- a/server/pom.xml +++ b/server/pom.xml @@ -41,6 +41,10 @@ org.apache.rocketmq rocketmq-client + + org.apache.rocketmq + rocketmq-acl + org.jboss.netty @@ -54,9 +58,5 @@ junit test - - com.aliyun.openservices - aliware-apache-rocketmq-cloud - diff --git a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java index fc27baca..0073c186 100644 --- a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java +++ b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java @@ -222,15 +222,23 @@ public class MQMessageUtils { } else { for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { int hashCode = database.hashCode(); + CanalEntry.EventType eventType = rowChange.getEventType(); + List columns = null; + if (eventType == CanalEntry.EventType.DELETE) { + columns = rowData.getBeforeColumnsList(); + } else { + columns = rowData.getAfterColumnsList(); + } + if (hashMode.autoPkHash) { // isEmpty use default pkNames - for (CanalEntry.Column column : rowData.getAfterColumnsList()) { + for (CanalEntry.Column column : columns) { if (column.getIsKey()) { hashCode = hashCode ^ column.getValue().hashCode(); } } } else { - for (CanalEntry.Column column : rowData.getAfterColumnsList()) { + for (CanalEntry.Column column : columns) { if (checkPkNamesHasContain(hashMode.pkNames, column.getName())) { hashCode = hashCode ^ column.getValue().hashCode(); } @@ -438,12 +446,14 @@ public class MQMessageUtils { int idx = 0; for (Map row : flatMessage.getData()) { int hashCode = database.hashCode(); - for (String pkName : pkNames) { - String value = row.get(pkName); - if (value == null) { - value = ""; + if (pkNames != null) { + for (String pkName : pkNames) { + String value = row.get(pkName); + if (value == null) { + value = ""; + } + hashCode = hashCode ^ value.hashCode(); } - hashCode = hashCode ^ value.hashCode(); } int pkHash = Math.abs(hashCode) % partitionsNum; @@ -463,6 +473,7 @@ public class MQMessageUtils { flatMessageTmp.setMysqlType(flatMessage.getMysqlType()); flatMessageTmp.setEs(flatMessage.getEs()); flatMessageTmp.setTs(flatMessage.getTs()); + flatMessageTmp.setPkNames(flatMessage.getPkNames()); } List> data = flatMessageTmp.getData(); if (data == null) { diff --git a/server/src/main/java/com/alibaba/otter/canal/common/MQProperties.java b/server/src/main/java/com/alibaba/otter/canal/common/MQProperties.java index 4ee28c5e..41e41684 100644 --- a/server/src/main/java/com/alibaba/otter/canal/common/MQProperties.java +++ b/server/src/main/java/com/alibaba/otter/canal/common/MQProperties.java @@ -3,7 +3,7 @@ package com.alibaba.otter.canal.common; import java.util.Properties; /** - * kafka 配置项 + * MQ 配置项 * * @author machengyuan 2018-6-11 下午05:30:49 * @version 1.0.0 @@ -27,6 +27,13 @@ public class MQProperties { private String aliyunSecretKey = ""; private boolean transaction = false; // 是否开启事务 private Properties properties = new Properties(); + private boolean enableMessageTrace = false; + private String accessChannel = null; + private String customizedTraceTopic = null; + private String namespace = ""; + private boolean kerberosEnable = false; //kafka集群是否启动Kerberos认证 + private String kerberosKrb5FilePath = ""; //启动Kerberos认证时配置为krb5.conf文件的路径 + private String kerberosJaasFilePath = ""; //启动Kerberos认证时配置为jaas.conf文件的路径 public static class CanalDestination { @@ -221,4 +228,89 @@ public class MQProperties { public void setProperties(Properties properties) { this.properties = properties; } + + public boolean isEnableMessageTrace() { + return enableMessageTrace; + } + + public void setEnableMessageTrace(boolean enableMessageTrace) { + this.enableMessageTrace = enableMessageTrace; + } + + public String getAccessChannel() { + return accessChannel; + } + + public void setAccessChannel(String accessChannel) { + this.accessChannel = accessChannel; + } + + public String getCustomizedTraceTopic() { + return customizedTraceTopic; + } + + public void setCustomizedTraceTopic(String customizedTraceTopic) { + this.customizedTraceTopic = customizedTraceTopic; + } + + public String getNamespace() { + return namespace; + } + + public void setNamespace(String namespace) { + this.namespace = namespace; + } + + public boolean isKerberosEnable() { + return kerberosEnable; + } + + public void setKerberosEnable(boolean kerberosEnable) { + this.kerberosEnable = kerberosEnable; + } + + public String getKerberosKrb5FilePath() { + return kerberosKrb5FilePath; + } + + public void setKerberosKrb5FilePath(String kerberosKrb5FilePath) { + this.kerberosKrb5FilePath = kerberosKrb5FilePath; + } + + public String getKerberosJaasFilePath() { + return kerberosJaasFilePath; + } + + public void setKerberosJaasFilePath(String kerberosJaasFilePath) { + this.kerberosJaasFilePath = kerberosJaasFilePath; + } + + @Override public String toString() { + return "MQProperties{" + + "servers='" + servers + '\'' + + ", retries=" + retries + + ", batchSize=" + batchSize + + ", lingerMs=" + lingerMs + + ", maxRequestSize=" + maxRequestSize + + ", bufferMemory=" + bufferMemory + + ", filterTransactionEntry=" + filterTransactionEntry + + ", producerGroup='" + producerGroup + '\'' + + ", canalBatchSize=" + canalBatchSize + + ", canalGetTimeout=" + canalGetTimeout + + ", flatMessage=" + flatMessage + + ", compressionType='" + compressionType + '\'' + + ", acks='" + acks + '\'' + + ", aliyunAccessKey='" + aliyunAccessKey + '\'' + + ", aliyunSecretKey='" + aliyunSecretKey + '\'' + + ", transaction=" + transaction + + ", properties=" + properties + + ", enableMessageTrace=" + enableMessageTrace + + ", accessChannel='" + accessChannel + '\'' + + ", customizedTraceTopic='" + customizedTraceTopic + '\'' + + ", namespace='" + namespace + '\'' + + ", kerberosEnable='" + kerberosEnable + '\'' + + ", kerberosKrb5FilePath='" + kerberosKrb5FilePath + '\'' + + ", kerberosJaasFilePath='" + kerberosJaasFilePath + '\'' + + '}'; + } } diff --git a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java index 89690a29..99f2010a 100644 --- a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java @@ -1,19 +1,5 @@ package com.alibaba.otter.canal.kafka; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.Properties; -import java.util.concurrent.ExecutionException; - -import org.apache.commons.lang.StringUtils; -import org.apache.kafka.clients.producer.KafkaProducer; -import org.apache.kafka.clients.producer.Producer; -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.common.serialization.StringSerializer; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.serializer.SerializerFeature; import com.alibaba.otter.canal.common.MQMessageUtils; @@ -21,6 +7,20 @@ import com.alibaba.otter.canal.common.MQProperties; import com.alibaba.otter.canal.protocol.FlatMessage; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.spi.CanalMQProducer; +import org.apache.commons.lang.StringUtils; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringSerializer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.File; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.concurrent.ExecutionException; /** * kafka producer 主操作类 @@ -59,6 +59,26 @@ public class CanalKafkaProducer implements CanalMQProducer { } else { properties.put("retries", kafkaProperties.getRetries()); } + + if (kafkaProperties.isKerberosEnable()){ + File krb5File = new File(kafkaProperties.getKerberosKrb5FilePath()); + File jaasFile = new File(kafkaProperties.getKerberosJaasFilePath()); + if(krb5File.exists() && jaasFile.exists()){ + //配置kerberos认证,需要使用绝对路径 + System.setProperty("java.security.krb5.conf", + krb5File.getAbsolutePath()); + System.setProperty("java.security.auth.login.config", + jaasFile.getAbsolutePath()); + System.setProperty("javax.security.auth.useSubjectCredsOnly", "false"); + properties.put("security.protocol", "SASL_PLAINTEXT"); + properties.put("sasl.kerberos.service.name", "kafka"); + }else{ + String errorMsg = "ERROR # The kafka kerberos configuration file does not exist! please check it"; + logger.error(errorMsg); + throw new RuntimeException(errorMsg); + } + } + if (!kafkaProperties.getFlatMessage()) { properties.put("value.serializer", MessageSerializer.class.getName()); producer = new KafkaProducer(properties); @@ -102,10 +122,10 @@ public class CanalKafkaProducer implements CanalMQProducer { producerTmp = producer2; } - if (kafkaProperties.getTransaction()) { - producerTmp.beginTransaction(); - } try { + if (kafkaProperties.getTransaction()) { + producerTmp.beginTransaction(); + } if (!StringUtils.isEmpty(canalDestination.getDynamicTopic())) { // 动态topic Map messageMap = MQMessageUtils.messageTopics(message, @@ -130,7 +150,11 @@ public class CanalKafkaProducer implements CanalMQProducer { } catch (Throwable e) { logger.error(e.getMessage(), e); if (kafkaProperties.getTransaction()) { - producerTmp.abortTransaction(); + try { + producerTmp.abortTransaction(); + } catch (Exception e1) { + logger.error(e1.getMessage(), e1); + } } callback.rollback(); } diff --git a/server/src/main/java/com/alibaba/otter/canal/rocketmq/CanalRocketMQProducer.java b/server/src/main/java/com/alibaba/otter/canal/rocketmq/CanalRocketMQProducer.java index aaa658e2..e16ab3d3 100644 --- a/server/src/main/java/com/alibaba/otter/canal/rocketmq/CanalRocketMQProducer.java +++ b/server/src/main/java/com/alibaba/otter/canal/rocketmq/CanalRocketMQProducer.java @@ -4,17 +4,20 @@ import java.util.List; import java.util.Map; import org.apache.commons.lang.StringUtils; +import org.apache.rocketmq.acl.common.AclClientRPCHook; +import org.apache.rocketmq.acl.common.SessionCredentials; +import org.apache.rocketmq.client.AccessChannel; import org.apache.rocketmq.client.exception.MQBrokerException; import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.MessageQueueSelector; +import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.remoting.RPCHook; import org.apache.rocketmq.remoting.exception.RemotingException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.serializer.SerializerFeature; import com.alibaba.otter.canal.common.CanalMessageSerializer; @@ -23,14 +26,13 @@ import com.alibaba.otter.canal.common.MQProperties; import com.alibaba.otter.canal.protocol.FlatMessage; import com.alibaba.otter.canal.server.exception.CanalServerException; import com.alibaba.otter.canal.spi.CanalMQProducer; -import com.aliyun.openservices.apache.api.impl.authority.SessionCredentials; -import com.aliyun.openservices.apache.api.impl.rocketmq.ClientRPCHook; public class CanalRocketMQProducer implements CanalMQProducer { private static final Logger logger = LoggerFactory.getLogger(CanalRocketMQProducer.class); private DefaultMQProducer defaultMQProducer; private MQProperties mqProperties; + private static final String CLOUD_ACCESS_CHANNEL = "cloud"; @Override public void init(MQProperties rocketMQProperties) { @@ -41,10 +43,16 @@ public class CanalRocketMQProducer implements CanalMQProducer { SessionCredentials sessionCredentials = new SessionCredentials(); sessionCredentials.setAccessKey(rocketMQProperties.getAliyunAccessKey()); sessionCredentials.setSecretKey(rocketMQProperties.getAliyunSecretKey()); - rpcHook = new ClientRPCHook(sessionCredentials); + rpcHook = new AclClientRPCHook(sessionCredentials); } - defaultMQProducer = new DefaultMQProducer(rocketMQProperties.getProducerGroup(), rpcHook); + defaultMQProducer = new DefaultMQProducer(rocketMQProperties.getProducerGroup(), rpcHook, mqProperties.isEnableMessageTrace(), mqProperties.getCustomizedTraceTopic()); + if (CLOUD_ACCESS_CHANNEL.equals(rocketMQProperties.getAccessChannel())){ + defaultMQProducer.setAccessChannel(AccessChannel.CLOUD); + } + if (!StringUtils.isEmpty(mqProperties.getNamespace())){ + defaultMQProducer.setNamespace(mqProperties.getNamespace()); + } defaultMQProducer.setNamesrvAddr(rocketMQProperties.getServers()); defaultMQProducer.setRetryTimesWhenSendFailed(rocketMQProperties.getRetries()); defaultMQProducer.setVipChannelEnabled(false); @@ -80,6 +88,22 @@ public class CanalRocketMQProducer implements CanalMQProducer { } } + private void sendMessage(Message message, int partition) throws Exception{ + SendResult sendResult = this.defaultMQProducer.send(message, new MessageQueueSelector() { + @Override + public MessageQueue select(List mqs, Message msg, Object arg) { + if (partition > mqs.size()) { + return mqs.get(partition % mqs.size()); + } else { + return mqs.get(partition); + } + } + }, null); + if (logger.isDebugEnabled()) { + logger.debug("Send Message Result: {}", sendResult); + } + } + public void send(final MQProperties.CanalDestination destination, String topicName, com.alibaba.otter.canal.protocol.Message data) throws Exception { if (!mqProperties.getFlatMessage()) { @@ -102,17 +126,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { Message message = new Message(topicName, CanalMessageSerializer.serializer(dataPartition, mqProperties.isFilterTransactionEntry())); - this.defaultMQProducer.send(message, new MessageQueueSelector() { - - @Override - public MessageQueue select(List mqs, Message msg, Object arg) { - if (index > mqs.size()) { - return mqs.get(index % mqs.size()); - } else { - return mqs.get(index); - } - } - }, null); + sendMessage(message, index); } catch (Exception e) { logger.error("send flat message to hashed partition error", e); throw e; @@ -129,17 +143,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { destination.getCanalDestination(), partition); } - this.defaultMQProducer.send(message, new MessageQueueSelector() { - - @Override - public MessageQueue select(List mqs, Message msg, Object arg) { - if (partition > mqs.size()) { - return mqs.get(partition % mqs.size()); - } else { - return mqs.get(partition); - } - } - }, null); + sendMessage(message, partition); } } catch (MQClientException | RemotingException | MQBrokerException | InterruptedException e) { logger.error("Send message error!", e); @@ -166,17 +170,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { try { Message message = new Message(topicName, JSON.toJSONString(flatMessagePart, SerializerFeature.WriteMapNullValue).getBytes()); - this.defaultMQProducer.send(message, new MessageQueueSelector() { - - @Override - public MessageQueue select(List mqs, Message msg, Object arg) { - if (index > mqs.size()) { - return mqs.get(index % mqs.size()); - } else { - return mqs.get(index); - } - } - }, null); + sendMessage(message, index); } catch (Exception e) { logger.error("send flat message to hashed partition error", e); throw e; @@ -194,17 +188,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { } Message message = new Message(topicName, JSON.toJSONString(flatMessage, SerializerFeature.WriteMapNullValue).getBytes()); - this.defaultMQProducer.send(message, new MessageQueueSelector() { - - @Override - public MessageQueue select(List mqs, Message msg, Object arg) { - if (partition > mqs.size()) { - return mqs.get(partition % mqs.size()); - } else { - return mqs.get(partition); - } - } - }, null); + sendMessage(message, partition); } catch (Exception e) { logger.error("send flat message to fixed partition error", e); throw e;