From beed10901d210d5fc6ad45cebeca2c0df3f27285 Mon Sep 17 00:00:00 2001 From: wuwo Date: Mon, 5 Nov 2018 16:22:15 +0800 Subject: [PATCH 1/5] destory TableMetaTSDB when stop EventParser --- .../main/resources/spring/tsdb/h2-tsdb.xml | 2 +- .../main/resources/spring/tsdb/mysql-tsdb.xml | 2 +- .../inbound/mysql/tsdb/DatabaseTableMeta.java | 22 ++++++++++++++++++- .../inbound/mysql/tsdb/MemoryTableMeta.java | 5 +++++ .../inbound/mysql/tsdb/TableMetaTSDB.java | 5 +++++ parse/src/test/resources/tsdb/h2-tsdb.xml | 2 +- parse/src/test/resources/tsdb/mysql-tsdb.xml | 2 +- 7 files changed, 35 insertions(+), 5 deletions(-) diff --git a/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml b/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml index beb61236..c05b5200 100644 --- a/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml +++ b/deployer/src/main/resources/spring/tsdb/h2-tsdb.xml @@ -23,7 +23,7 @@ - + diff --git a/deployer/src/main/resources/spring/tsdb/mysql-tsdb.xml b/deployer/src/main/resources/spring/tsdb/mysql-tsdb.xml index 58cab638..87f407b3 100644 --- a/deployer/src/main/resources/spring/tsdb/mysql-tsdb.xml +++ b/deployer/src/main/resources/spring/tsdb/mysql-tsdb.xml @@ -23,7 +23,7 @@ - + diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java index b0dd89ce..73a5ead8 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java @@ -72,7 +72,7 @@ public class DatabaseTableMeta implements TableMetaTSDB { @Override public Thread newThread(Runnable r) { - Thread thread = new Thread(r, "[scheduler-table-meta-snapshot]"); + Thread thread = new Thread(r, String.format("[scheduler-table-meta-snapshot-%s]", destination)); thread.setDaemon(true); return thread; } @@ -105,6 +105,26 @@ public class DatabaseTableMeta implements TableMetaTSDB { } return true; } + + @Override + public void destory() { + if (memoryTableMeta != null) { + memoryTableMeta.destory(); + } + + if (connection != null) { + try { + connection.disconnect(); + } catch (IOException e) { + logger.error("ERROR # disconnect meta connection for address:{}", connection.getConnector() + .getAddress(), e); + } + } + + if (scheduler != null) { + scheduler.shutdown(); + } + } @Override public TableMeta find(String schema, String table) { 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 904b55a9..c53e1a08 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 @@ -58,6 +58,11 @@ public class MemoryTableMeta implements TableMetaTSDB { public boolean init(String destination) { return true; } + + @Override + public void destory() { + tableMetas.clear(); + } public boolean apply(EntryPosition position, String schema, String ddl, String extra) { tableMetas.clear(); diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDB.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDB.java index d712adbe..f968ae12 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDB.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDB.java @@ -18,6 +18,11 @@ public interface TableMetaTSDB { */ public boolean init(String destination); + /** + * 销毁资源 + */ + public void destory(); + /** * 获取当前的表结构 */ diff --git a/parse/src/test/resources/tsdb/h2-tsdb.xml b/parse/src/test/resources/tsdb/h2-tsdb.xml index c0a3d385..dcb6bc47 100644 --- a/parse/src/test/resources/tsdb/h2-tsdb.xml +++ b/parse/src/test/resources/tsdb/h2-tsdb.xml @@ -7,7 +7,7 @@ default-autowire="byName"> - + diff --git a/parse/src/test/resources/tsdb/mysql-tsdb.xml b/parse/src/test/resources/tsdb/mysql-tsdb.xml index 099c034d..44f97203 100644 --- a/parse/src/test/resources/tsdb/mysql-tsdb.xml +++ b/parse/src/test/resources/tsdb/mysql-tsdb.xml @@ -7,7 +7,7 @@ default-autowire="byName"> - + From e608f09b53e1d0d0964fda73693320f3f66148dc Mon Sep 17 00:00:00 2001 From: wuwo Date: Mon, 5 Nov 2018 16:37:42 +0800 Subject: [PATCH 2/5] improve DatabaseTableMeta --- .../inbound/mysql/tsdb/DatabaseTableMeta.java | 28 +++++++++---------- 1 file changed, 14 insertions(+), 14 deletions(-) diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java index 73a5ead8..8896e2f5 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DatabaseTableMeta.java @@ -7,6 +7,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; import java.util.regex.Pattern; @@ -48,18 +49,26 @@ public class DatabaseTableMeta implements TableMetaTSDB { private static Logger logger = LoggerFactory.getLogger(DatabaseTableMeta.class); private static Pattern pattern = Pattern.compile("Duplicate entry '.*' for key '*'"); private static Pattern h2Pattern = Pattern.compile("Unique index or primary key violation"); + private static ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(new ThreadFactory() { + @Override + public Thread newThread(Runnable r) { + Thread thread = new Thread(r, "[scheduler-table-meta-snapshot]"); + thread.setDaemon(true); + return thread; + } + }); private String destination; private MemoryTableMeta memoryTableMeta; private MysqlConnection connection; // 查询meta信息的链接 private CanalEventFilter filter; private CanalEventFilter blackFilter; private EntryPosition lastPosition; - private ScheduledExecutorService scheduler; private MetaHistoryDAO metaHistoryDAO; private MetaSnapshotDAO metaSnapshotDAO; private int snapshotInterval = 24; private int snapshotExpire = 360; - + private ScheduledFuture scheduleSnapshotFuture; + public DatabaseTableMeta(){ } @@ -68,19 +77,10 @@ public class DatabaseTableMeta implements TableMetaTSDB { public boolean init(final String destination) { this.destination = destination; this.memoryTableMeta = new MemoryTableMeta(); - this.scheduler = Executors.newSingleThreadScheduledExecutor(new ThreadFactory() { - - @Override - public Thread newThread(Runnable r) { - Thread thread = new Thread(r, String.format("[scheduler-table-meta-snapshot-%s]", destination)); - thread.setDaemon(true); - return thread; - } - }); // 24小时生成一份snapshot if (snapshotInterval > 0) { - scheduler.scheduleWithFixedDelay(new Runnable() { + scheduleSnapshotFuture = scheduler.scheduleWithFixedDelay(new Runnable() { @Override public void run() { @@ -121,8 +121,8 @@ public class DatabaseTableMeta implements TableMetaTSDB { } } - if (scheduler != null) { - scheduler.shutdown(); + if (scheduleSnapshotFuture != null) { + scheduleSnapshotFuture.cancel(false); } } From 359aa1a0740d2b47d6ebdc70377b9fbdfa42c5df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=83=E9=94=8B?= Date: Mon, 5 Nov 2018 16:53:15 +0800 Subject: [PATCH 3/5] fixed compiler failed --- .../dbsync/binlog/BaseLogFetcherTest.java | 4 +- pom.xml | 57 +------------------ 2 files changed, 4 insertions(+), 57 deletions(-) diff --git a/dbsync/src/test/java/com/taobao/tddl/dbsync/binlog/BaseLogFetcherTest.java b/dbsync/src/test/java/com/taobao/tddl/dbsync/binlog/BaseLogFetcherTest.java index d0c6746d..27ecf3db 100644 --- a/dbsync/src/test/java/com/taobao/tddl/dbsync/binlog/BaseLogFetcherTest.java +++ b/dbsync/src/test/java/com/taobao/tddl/dbsync/binlog/BaseLogFetcherTest.java @@ -68,7 +68,7 @@ public class BaseLogFetcherTest { // update需要处理before/after System.out.println("-------> before"); parseOneRow(event, buffer, columns, false); - if (!buffer.nextOneRow(changeColumns)) { + if (!buffer.nextOneRow(changeColumns, true)) { break; } System.out.println("-------> after"); @@ -97,7 +97,7 @@ public class BaseLogFetcherTest { } ColumnInfo info = columnInfo[i]; - buffer.nextValue(info.type, info.meta); + buffer.nextValue(null , i ,info.type, info.meta); if (buffer.isNull()) { // diff --git a/pom.xml b/pom.xml index 5cb51eac..0bb4bc11 100644 --- a/pom.xml +++ b/pom.xml @@ -371,6 +371,7 @@ + + --> src/main/java src/test/java From d8e411462ec26ee8644445280fd88734b53a09d0 Mon Sep 17 00:00:00 2001 From: winger Date: Mon, 5 Nov 2018 18:38:44 +0800 Subject: [PATCH 4/5] =?UTF-8?q?=E5=A2=9E=E5=8A=A0alter=20rename=20?= =?UTF-8?q?=E7=9A=84=E7=B1=BB=E5=9E=8B=E4=B8=BArename?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../otter/canal/parse/inbound/mysql/ddl/DruidDdlParser.java | 1 + 1 file changed, 1 insertion(+) diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/ddl/DruidDdlParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/ddl/DruidDdlParser.java index d82e5eff..68f1ff5a 100644 --- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/ddl/DruidDdlParser.java +++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/ddl/DruidDdlParser.java @@ -71,6 +71,7 @@ public class DruidDdlParser { DdlResult ddlResult = new DdlResult(); processName(ddlResult, schmeaName, alterTable.getName(), true); processName(ddlResult, schmeaName, ((SQLAlterTableRename) item).getToName(), false); + ddlResult.setType(EventType.RENAME); ddlResults.add(ddlResult); } else if (item instanceof SQLAlterTableAddIndex) { DdlResult ddlResult = new DdlResult(); From edb5f294ef45174fbd9b7392e16ca17c9a15ebef Mon Sep 17 00:00:00 2001 From: jiacheo Date: Mon, 5 Nov 2018 22:00:56 +0800 Subject: [PATCH 5/5] =?UTF-8?q?fix=EF=BC=9Arollback=E4=B9=8B=E5=90=8E?= =?UTF-8?q?=E7=BB=A7=E7=BB=ADcommit=EF=BC=8C=E4=BC=9A=E4=BA=A7=E7=94=9F=20?= =?UTF-8?q?CanalServerException:=20ack=20error=20,=20clientId:1001=20batch?= =?UTF-8?q?Id:12=20is=20not=20exist=20,=20please=20check=20=E7=9A=84?= =?UTF-8?q?=E9=94=99=E8=AF=AF=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../otter/canal/kafka/CanalKafkaProducer.java | 23 +++++++++++-------- .../canal/rocketmq/CanalRocketMQProducer.java | 18 ++++++++------- 2 files changed, 23 insertions(+), 18 deletions(-) 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 654a4b6b..8576835f 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,21 +1,20 @@ package com.alibaba.otter.canal.kafka; -import java.util.List; -import java.util.Properties; - -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.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.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.util.List; +import java.util.Properties; /** * kafka producer 主操作类 @@ -93,6 +92,7 @@ public class CanalKafkaProducer implements CanalMQProducer { logger.error(e.getMessage(), e); // producer.abortTransaction(); callback.rollback(); + return; } } else { // 发送扁平数据json @@ -110,6 +110,7 @@ public class CanalKafkaProducer implements CanalMQProducer { logger.error(e.getMessage(), e); // producer.abortTransaction(); callback.rollback(); + return; } } else { if (canalDestination.getPartitionHash() != null @@ -131,6 +132,7 @@ public class CanalKafkaProducer implements CanalMQProducer { logger.error(e.getMessage(), e); // producer.abortTransaction(); callback.rollback(); + return; } } } @@ -145,6 +147,7 @@ public class CanalKafkaProducer implements CanalMQProducer { logger.error(e.getMessage(), e); // producer.abortTransaction(); callback.rollback(); + return; } } } 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 6173cebd..d8bd45e4 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 @@ -1,7 +1,11 @@ package com.alibaba.otter.canal.rocketmq; -import java.util.List; - +import com.alibaba.fastjson.JSON; +import com.alibaba.otter.canal.common.CanalMessageSerializer; +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 org.apache.rocketmq.client.exception.MQBrokerException; import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.producer.DefaultMQProducer; @@ -12,12 +16,7 @@ import org.apache.rocketmq.remoting.exception.RemotingException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.alibaba.fastjson.JSON; -import com.alibaba.otter.canal.common.CanalMessageSerializer; -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 java.util.List; public class CanalRocketMQProducer implements CanalMQProducer { @@ -67,6 +66,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { } catch (MQClientException | RemotingException | MQBrokerException | InterruptedException e) { logger.error("Send message error!", e); callback.rollback(); + return; } } else { List flatMessages = FlatMessage.messageConverter(data); @@ -90,6 +90,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { } catch (Exception e) { logger.error("send flat message to fixed partition error", e); callback.rollback(); + return; } } else { if (destination.getPartitionHash() != null && !destination.getPartitionHash().isEmpty()) { @@ -124,6 +125,7 @@ public class CanalRocketMQProducer implements CanalMQProducer { } catch (Exception e) { logger.error("send flat message to hashed partition error", e); callback.rollback(); + return; } } }