Compare commits

...
Author SHA1 Message Date
七锋 c340bd6565 upgrade fastsql 159 2018-03-16 15:10:02 +08:00
七锋 511734c3d8 fixed compiler error 2018-03-15 20:21:42 +08:00
七锋 af002f4e24 fixed bio socket 2018-03-15 20:20:29 +08:00
七锋 28ad12d0ea upgrade fastsql 2.0.0_preview_151 2018-03-15 15:48:16 +08:00
agapple 6229a8593f Merge pull request #555 from whaon/patch-1
renew an InetSocketAddress to resolve address again
2018-03-15 14:33:05 +08:00
whaon 256056e14f renew an InetSocketAddress to resolve address again
resolve [issue427](https://github.com/alibaba/otter/issues/427)
2018-03-14 09:55:20 +08:00
agapple 499be8bd8c upgrade druid to fastsql 2018-03-13 16:28:08 +08:00
agapple cd87ef00cc fixed issue #553 , merge code 2018-03-13 10:00:11 +08:00
agapple 6b7109c66b fixed issue #539 , support bio&netty 2018-03-12 23:21:43 +08:00
七锋 b90c9fac41 fixed druid version 1.1.9 2018-03-12 17:14:12 +08:00
agapple 88c79544f0 Merge pull request #532 from LavaCoref/hotfix/lava/1_json
fix(dbsync): json column of zero length has no value, value parsing s…
2018-03-12 16:37:50 +08:00
agapple fddaff1bd9 Merge pull request #536 from Wu-Jianqiang/master
主要修复 SocketChannel 默认缓存大小(1MB)不够用时需自动扩充,否则将因缓存空间不足而造成I/O超时假象
2018-03-12 16:36:07 +08:00
七锋 847ebf9bf3 fixed issue #535 , use show create table to get tableMeta 2018-03-12 16:27:58 +08:00
Wu-Jianqiang 3e7e7b7f6e 修改:变量命名和super调用 2018-02-26 16:37:21 +08:00
Wu-Jianqiang 06e7f44114 修复:MASTER_HEARTBEAT_PERIOD 的计量单位是秒(可以包含3位小数表示毫秒),但不是纳秒。 2018-02-26 16:31:58 +08:00
Wu-Jianqiang 8ad603fea4 修复:SocketChannel 默认缓存大小(1MB)不够用时需自动扩充,否则将因缓存空间不足而造成I/O超时假象 2018-02-26 16:26:01 +08:00
WU Jianqiang 863a61d626 Merge pull request #2 from alibaba/master
merge
2018-02-23 23:17:46 +08:00
lvtianyu 7b38542183 fix(dbsync): json column of zero length has no value, value parsing should be skipped 2018-02-23 10:41:49 +08:00
agapple 0af789b1c7 Merge pull request #529 from yakirChen/master
TimelineTransactionBarrier import CanalSinkException; Maven Compile 的一些 Warning
2018-02-12 11:43:57 +08:00
agapple 2d49b709ca fixed socket wait period 10ms 2018-02-12 11:43:03 +08:00
agapple 826a45c105 Merge pull request #528 from lcybo/timeout
Use mysql master heartbeat to detect phycial tcp connection failure.
2018-02-12 11:41:55 +08:00
yakirChen 03452dcdda Merge branch 'master' of github.com:yakirChen/canal 2018-02-12 11:29:16 +08:00
yakirChen 514876f60f [x] Fix Compile Error : TimelineTransactionBarrier.java:[70,27] 找不到符号: 类 CanalSinkException
[x] Fix Maven Wranning: The expression ${pom.version} is deprecated. Please use ${project.version} instead.
[x] Fix MavenWranning : 'build.plugins.plugin.version' for org.apache.maven.plugins:maven-jar-plugin is missing.
2018-02-12 11:28:42 +08:00
yakirChen 89994c52f9 [x] Fix Compile Error : TimelineTransactionBarrier.java:[70,27] 找不到符号: 类 CanalSinkException
[x] Fix Maven Wranning: The expression ${pom.version} is deprecated. Please use ${project.version} instead.
[x] Fix MavenWranning : 'build.plugins.plugin.version' for org.apache.maven.plugins:maven-jar-plugin is missing.
2018-02-12 11:17:54 +08:00
agapple edc2fbe1dc fixed issue #507 2018-02-12 11:15:43 +08:00
agapple cfa90af33c fixed unique and upgrade druid 2018-02-12 11:01:53 +08:00
Chuanyi Li 32d4174083 Merge branch 'timeout' of https://github.com/lcybo/canal into timeout 2018-02-11 00:16:07 +08:00
Chuanyi Li bb0f3b3a7e [MI] Use mysql master heartbeat to detect phycial tcp connection
failure.
2018-02-11 00:10:16 +08:00
Chuanyi Li dc0010823f [MI] Use mysql master heartbeat to detect phycial tcp connection
failure.
2018-02-10 23:42:06 +08:00
agapple b1431f8827 Merge pull request #521 from lcybo/monitor
修改InstanceConfig未生成或已销毁时造成ServerRunningMonitor running状态与实际不一致的问题
2018-02-10 11:22:48 +08:00
agapple b9e3629459 Merge pull request #522 from lulu2panpan/lulu2panpan-patch-1
修复死锁bug
2018-02-10 11:18:22 +08:00
agapple 31ed4efe3a Merge pull request #497 from lcybo/master
通信模块log小优化
2018-02-10 11:15:44 +08:00
Chuanyi Li 6da1507405 Rollback running state and znode when start failed. 2018-02-06 20:54:33 +08:00
WU Jianqiang 0b82ee8d10 Merge pull request #1 from alibaba/master
pull & merge
2018-02-02 16:13:21 +08:00
Chuanyi Li 00890f87e5 Communication log optimization. 2018-01-18 21:11:46 +08:00
agapple 91ac0e30aa fixed issue #483 , support show slave hosts 2018-01-10 14:35:26 +08:00
agapple 9316a3df6d fixed issue #482 , retry for getTableMetaByDB 2018-01-10 13:49:25 +08:00
agapple 2d31e6bf91 Merge pull request #487 from jasonhx140/master
dump connection is disconnected
2018-01-10 13:41:30 +08:00
jason huang a965d2cd04 dump connection is disconnected
dump connection is disconnected since SocketChannel.cache is overflow.
2018-01-09 18:25:51 +08:00
jason huang 93be9d1909 dump connection disconnected
dump connection will be disconnected since SocketChannel.cache is overflow.
2018-01-09 18:21:14 +08:00
jason huang 7cf9375ac6 Merge pull request #2 from alibaba/master
merge new release
2018-01-09 18:10:19 +08:00
agapple af289cd9fb fixed compiler error 2018-01-05 19:09:46 +08:00
agapple d1f1c0716b Merge pull request #473 from jasonhx140/master
connect mysql server failed when reconnect mysql server frequently
2018-01-05 18:55:18 +08:00
jason huang 0c7aed197a cannot build mysql connection.
when invoke connect(), it will return a ChannelFuture instance, and then the "addListener" method is invoked by biz thread, however, "channelRead" method in BusinessHandler is executed by netty I/O thread.

ChannelFutureListener -> operationComplete() cannot guarantee that which is executed before than BusinessHandler -> channelRead() method.

reproduce it :
set below jvm parameter:
-Dio.netty.eventLoopThreads=1

It will caused new connection is timeout when AbstractEventParser -> start() -> parseThread ->erosaConnection.reconnect() is invoked.
2018-01-02 23:29:15 +08:00
jason huang 325d16bc66 Merge pull request #1 from alibaba/master
merge latest version
2018-01-02 23:03:52 +08:00
agapple 3104b4cc24 fixed issue #449 , optimize findTransactionBeginPosition 2017-12-27 17:34:46 +08:00
agapple 9595e6abbd fixed semi sync format & refactor 2017-12-27 16:53:09 +08:00
agapple c6f4de9f88 Merge pull request #455 from lucifax301/master
mysql semi support and mariadb gtid parse
2017-12-27 16:36:29 +08:00
lucifax301 d66839c965 mysql semi support and mariadb gtid parse 2017-12-18 14:10:53 +08:00
agapple ac773009f8 Merge pull request #447 from zhmz1326/master
解决destination无限连接等待的bug
2017-12-13 10:31:48 +08:00
zhengmingzhi 8b03f172ed 解决destination无限连接等待的bug 2017-12-12 19:17:53 +08:00
agapple 5e206e4f7e fixed issue #440 , optimize find position 2017-12-10 23:05:28 +08:00
agapple 448fc8bcde fixed issue #440 2017-12-10 22:27:14 +08:00
agapple c60b04424b fixed issue #440 , optime find position 2017-12-10 22:24:24 +08:00
agapple 5aed505429 fixed issue #440, enableTsdb 2017-12-08 23:21:40 +08:00
agapple 912f5db2c7 fixed issue #440, for null position 2017-12-08 23:20:15 +08:00
agapple 8b57b9a0a5 fixed issue #442 , compareTableMetaDbAndMemory 2017-12-08 22:54:22 +08:00
agapple d1c83ea7bb fixed fastjson autotype 2017-12-05 19:46:49 +08:00
agapple 7d485efd71 Update README.md 2017-12-05 17:01:15 +08:00
agapple 20942533e6 [maven-release-plugin] prepare for next development iteration 2017-12-04 14:43:29 +08:00
agapple 6b2d80a145 [maven-release-plugin] prepare release canal-1.0.25 2017-12-04 14:43:16 +08:00
agapple 01ef15223c fixed issue #415, ignore flush privileges 2017-12-04 13:19:58 +08:00
agapple da3290c381 Merge pull request #424 from lcybo/master
1.0.23里面修复的KILL DUMP exception问题在(#334)里面被覆盖回去了
2017-12-03 21:36:33 -06:00
agapple a9284b1b39 fixed compiler error 2017-12-04 11:35:02 +08:00
Chuanyi Li 3409472eb9 kill CONNECTION exception fix revert. 2017-11-23 22:25:47 +08:00
agapple 5348994b4a fixed alisql heartbeat 2017-11-15 10:42:58 +08:00
agapple a3b9f6f1eb fixed NPE 2017-11-14 23:44:15 +08:00
agapple 2dc019709b fixed quota ddl 2017-11-14 23:31:46 +08:00
agapple c56f58deae upgrade druid to 1.1.5 2017-11-01 20:29:47 +08:00
lulu2panpan fa059804ce 修复死锁bug
当通过zk-cursor启动的时候,第一个Event类型是TransactionEnd,那么txState变为2的触发条件不仅仅有ddl和dcl,clear方法中isTransactionEnd的判断应该放到txState.intValue() == 2的后面,否则将导致死锁
2017-10-19 10:50:32 +08:00
55 changed files with 1374 additions and 476 deletions
+1 -1
View File
@@ -36,7 +36,7 @@
<li>slave重做中继日志中的事件,将改变反映它自己的数据。</li>
</ol>
<h3>canal的工作原理:</h3>
<p><img width="590" src="https://camo.githubusercontent.com/46c626b4cde399db43b2634a7911a04aecf273a0/687474703a2f2f646c2e69746579652e636f6d2f75706c6f61642f6174746163686d656e742f303038302f333130372f63383762363762612d333934632d333038362d393537372d3964623035626530346339352e6a7067" alt="" height="273">
<p><img width="590" src="http://dl.iteye.com/upload/attachment/0080/3107/c87b67ba-394c-3086-9577-9db05be04c95.jpg" alt="" height="273">
<p>原理相对比较简单:</p>
<ol>
<li>canal模拟mysql slave的交互协议,伪装自己为mysql slave,向mysql master发送dump协议</li>
+2 -2
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -69,7 +69,7 @@
<links>
<link>https://github.com/alibaba/canal</link>
</links>
<outputDirectory>${project.build.directory}/apidocs/apidocs/${pom.version}</outputDirectory>
<outputDirectory>${project.build.directory}/apidocs/apidocs/${project.version}</outputDirectory>
</configuration>
</plugin>
<plugin>
+2 -3
View File
@@ -1,10 +1,9 @@
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<artifactId>canal.common</artifactId>
@@ -10,6 +10,7 @@ import java.util.List;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.TypeReference;
import com.alibaba.fastjson.parser.ParserConfig;
import com.alibaba.fastjson.serializer.JSONSerializer;
import com.alibaba.fastjson.serializer.ObjectSerializer;
import com.alibaba.fastjson.serializer.PropertyFilter;
@@ -28,6 +29,7 @@ public class JsonUtils {
SerializeConfig.getGlobalInstance().put(InetAddress.class, InetAddressSerializer.instance);
SerializeConfig.getGlobalInstance().put(Inet4Address.class, InetAddressSerializer.instance);
SerializeConfig.getGlobalInstance().put(Inet6Address.class, InetAddressSerializer.instance);
ParserConfig.getGlobalInstance().setAutoTypeSupport(true);
}
public static <T> T unmarshalFromByte(byte[] bytes, Class<T> targetClass) {
@@ -4,6 +4,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import com.alibaba.otter.canal.common.CanalException;
import org.I0Itec.zkclient.IZkDataListener;
import org.I0Itec.zkclient.exception.ZkException;
import org.I0Itec.zkclient.exception.ZkInterruptedException;
@@ -91,18 +92,25 @@ public class ServerRunningMonitor extends AbstractCanalLifeCycle {
processStart();
}
public void start() {
public synchronized void start() {
super.start();
processStart();
if (zkClient != null) {
// 如果需要尽可能释放instance资源,不需要监听running节点,不然即使stop了这台机器,另一台机器立马会start
String path = ZookeeperPathUtils.getDestinationServerRunning(destination);
zkClient.subscribeDataChanges(path, dataListener);
try {
processStart();
if (zkClient != null) {
// 如果需要尽可能释放instance资源,不需要监听running节点,不然即使stop了这台机器,另一台机器立马会start
String path = ZookeeperPathUtils.getDestinationServerRunning(destination);
zkClient.subscribeDataChanges(path, dataListener);
initRunning();
} else {
processActiveEnter();// 没有zk,直接启动
initRunning();
} else {
processActiveEnter();// 没有zk,直接启动
}
} catch (Exception e) {
logger.error("start failed", e);
// 没有正常启动,重置一下状态,避免干扰下一次start
stop();
}
}
public void release() {
@@ -113,7 +121,7 @@ public class ServerRunningMonitor extends AbstractCanalLifeCycle {
}
}
public void stop() {
public synchronized void stop() {
super.stop();
if (zkClient != null) {
@@ -234,11 +242,7 @@ public class ServerRunningMonitor extends AbstractCanalLifeCycle {
private void processActiveEnter() {
if (listener != null) {
try {
listener.processActiveEnter();
} catch (Exception e) {
logger.error("processActiveEnter failed", e);
}
listener.processActiveEnter();
}
}
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -20,6 +20,7 @@ public class LogBuffer {
protected int origin, limit;
protected int position;
protected int semival;
protected LogBuffer(){
}
@@ -104,9 +104,14 @@ public final class LogDecoder {
try {
/* Decoding binary-log to event */
event = decode(buffer, header, context);
if (event != null) {
event.setSemival(buffer.semival);
}
} catch (IOException e) {
if (logger.isWarnEnabled()) logger.warn("Decoding " + LogEvent.getTypeName(header.getType())
+ " failed from: " + context.getLogPosition(), e);
if (logger.isWarnEnabled()) {
logger.warn("Decoding " + LogEvent.getTypeName(header.getType()) + " failed from: "
+ context.getLogPosition(), e);
}
throw e;
} finally {
buffer.limit(limit); /* Restore limit */
@@ -359,8 +359,28 @@ public abstract class LogEvent {
protected static final Log logger = LogFactory.getLog(LogEvent.class);
protected final LogHeader header;
/**
* mysql半同步semi标识
*
* <pre>
* 0不需要semi ack 给mysql
* 1需要semi ack给mysql
* </pre>
*/
protected int semival;
protected LogEvent(LogHeader header){
public int getSemival() {
return semival;
}
public void setSemival(int semival) {
this.semival = semival;
}
protected LogEvent(LogHeader header){
this.header = header;
}
@@ -969,17 +969,17 @@ public final class RowsLogBuffer {
default:
throw new IllegalArgumentException("!! Unknown JSON packlen = " + meta);
}
// len = buffer.getUint16();
// buffer.forward(meta - 4);
int position = buffer.position();
Json_Value jsonValue = JsonConversion.parse_value(buffer.getUint8(), buffer, len - 1, charsetName);
StringBuilder builder = new StringBuilder();
jsonValue.toJsonString(builder, charsetName);
value = builder.toString();
buffer.position(position + len);
// byte[] binary = new byte[len];
// buffer.fillBytes(binary, 0, len);
// value = binary;
if (0 == len) {
// fixed issue #1 by lava, json column of zero length has no value, value parsing should be skipped
value = "";
} else {
int position = buffer.position();
Json_Value jsonValue = JsonConversion.parse_value(buffer.getUint8(), buffer, len - 1, charsetName);
StringBuilder builder = new StringBuilder();
jsonValue.toJsonString(builder, charsetName);
value = builder.toString();
buffer.position(position + len);
}
javaType = Types.VARCHAR;
length = len;
break;
@@ -13,9 +13,30 @@ import com.taobao.tddl.dbsync.binlog.event.LogHeader;
*/
public class MariaGtidLogEvent extends IgnorableLogEvent {
private long gtid;
/**
* <pre>
* mariadb gtidlog event format
* uint<8> GTID sequence
* uint<4> Replication Domain ID
* uint<1> Flags
*
* if flag & FL_GROUP_COMMIT_ID
* uint<8> commit_id
* else
* uint<6> 0
* </pre>
*/
public MariaGtidLogEvent(LogHeader header, LogBuffer buffer, FormatDescriptionLogEvent descriptionEvent){
super(header, buffer, descriptionEvent);
gtid = buffer.getUlong64().longValue();
// do nothing , just ignore log event
}
public long getGtid() {
return gtid;
}
}
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -32,6 +32,8 @@ public class CanalConstants {
public static final String CANAL_DESTINATION_PROPERTY = ROOT + ".instance.destination";
public static final String CANAL_SOCKETCHANNEL = ROOT + "." + "socketChannel";
public static String getInstanceModeKey(String destination) {
return MessageFormat.format(INSTANCE_MODE_TEMPLATE, destination);
}
@@ -85,6 +85,12 @@ public class CanalController {
// 初始化instance config
initInstanceConfig(properties);
// init socketChannel
String socketChannel = getProperty(properties, CanalConstants.CANAL_SOCKETCHANNEL);
if (StringUtils.isNotEmpty(socketChannel)) {
System.setProperty(CanalConstants.CANAL_SOCKETCHANNEL, socketChannel);
}
// 准备canal server
cid = Long.valueOf(getProperty(properties, CanalConstants.CANAL_ID));
ip = getProperty(properties, CanalConstants.CANAL_IP);
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -44,11 +44,16 @@ public class MysqlConnector {
private long connectionId = -1;
private AtomicBoolean connected = new AtomicBoolean(false);
public static final int timeout = 5 * 1000; // 5s
public MysqlConnector(){
}
public MysqlConnector(InetSocketAddress address, String username, String password){
this.address = address;
String addr = address.getHostString();
int port = address.getPort();
this.address = new InetSocketAddress(addr, port);
this.username = username;
this.password = password;
}
@@ -101,7 +106,8 @@ public class MysqlConnector {
MysqlUpdateExecutor executor = new MysqlUpdateExecutor(connector);
executor.update("KILL CONNECTION " + connectionId);
} catch (Exception e) {
throw new IOException("KILL DUMP " + connectionId + " failure", e);
// 忽略具体异常
logger.info("KILL DUMP " + connectionId + " failure", e);
} finally {
if (connector != null) {
connector.disconnect();
@@ -144,8 +150,8 @@ public class MysqlConnector {
}
private void negotiate(SocketChannel channel) throws IOException {
HeaderPacket header = PacketManager.readHeader(channel, 4);
byte[] body = PacketManager.readBytes(channel, header.getPacketBodyLength());
HeaderPacket header = PacketManager.readHeader(channel, 4, timeout);
byte[] body = PacketManager.readBytes(channel, header.getPacketBodyLength(), timeout);
if (body[0] < 0) {// check field_count
if (body[0] == -1) {
ErrorPacket error = new ErrorPacket();
@@ -184,7 +190,7 @@ public class MysqlConnector {
header = null;
header = PacketManager.readHeader(channel, 4);
body = null;
body = PacketManager.readBytes(channel, header.getPacketBodyLength());
body = PacketManager.readBytes(channel, header.getPacketBodyLength(), timeout);
assert body != null;
if (body[0] < 0) {
if (body[0] == -1) {
@@ -333,4 +339,9 @@ public class MysqlConnector {
public void setConnTimeout(int connTimeout) {
this.connTimeout = connTimeout;
}
public String getPassword() {
return password;
}
}
@@ -0,0 +1,57 @@
package com.alibaba.otter.canal.parse.driver.mysql.packets.client;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import com.alibaba.otter.canal.parse.driver.mysql.packets.CommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.utils.ByteHelper;
/**
* COM_REGISTER_SLAVE
*
* @author zhibinliu
* @since 1.0.24
*/
public class RegisterSlaveCommandPacket extends CommandPacket {
public String reportHost;
public int reportPort;
public String reportUser;
public String reportPasswd;
public long serverId;
public RegisterSlaveCommandPacket(){
setCommand((byte) 0x15);
}
public void fromBytes(byte[] data) {
// bypass
}
public static byte[] toLH(int n) {
byte[] b = new byte[4];
b[0] = (byte) (n & 0xff);
b[1] = (byte) (n >> 8 & 0xff);
b[2] = (byte) (n >> 16 & 0xff);
b[3] = (byte) (n >> 24 & 0xff);
return b;
}
public byte[] toBytes() throws IOException {
ByteArrayOutputStream out = new ByteArrayOutputStream();
out.write(getCommand());
ByteHelper.writeUnsignedIntLittleEndian(serverId, out);
out.write((byte) reportHost.getBytes().length);
ByteHelper.writeFixedLengthBytesFromStart(reportHost.getBytes(), reportHost.getBytes().length, out);
out.write((byte) reportUser.getBytes().length);
ByteHelper.writeFixedLengthBytesFromStart(reportUser.getBytes(), reportUser.getBytes().length, out);
out.write((byte) reportPasswd.getBytes().length);
ByteHelper.writeFixedLengthBytesFromStart(reportPasswd.getBytes(), reportPasswd.getBytes().length, out);
ByteHelper.writeUnsignedShortLittleEndian(reportPort, out);
ByteHelper.writeUnsignedIntLittleEndian(0, out);// Fake
// rpl_recovery_rank
ByteHelper.writeUnsignedIntLittleEndian(0, out);// master id
return out.toByteArray();
}
}
@@ -0,0 +1,55 @@
package com.alibaba.otter.canal.parse.driver.mysql.packets.client;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import org.apache.commons.lang.StringUtils;
import com.alibaba.otter.canal.parse.driver.mysql.packets.CommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.utils.ByteHelper;
/**
* semi ack command
*
* @author amos_chen
*/
public class SemiAckCommandPacket extends CommandPacket {
public long binlogPosition;
public String binlogFileName;
public SemiAckCommandPacket(){
}
@Override
public void fromBytes(byte[] data) throws IOException {
}
/**
* <pre>
* Bytes Name
* --------------------------------------------------------
* Bytes Name
* ----- ----
* 1 semi mark
* 8 binlog position to start at (little endian)
* n binlog file name
*
* </pre>
*/
public byte[] toBytes() throws IOException {
ByteArrayOutputStream out = new ByteArrayOutputStream();
// 0 write semi mark
out.write(0xef);
// 1 write 8 bytes for position
ByteHelper.write8ByteUnsignedIntLittleEndian(binlogPosition, out);
// 2 write binlog filename
if (StringUtils.isNotEmpty(binlogFileName)) {
out.write(binlogFileName.getBytes());
}
return out.toByteArray();
}
}
@@ -0,0 +1,136 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.Socket;
import java.net.SocketAddress;
import java.net.SocketException;
import java.net.SocketTimeoutException;
import java.nio.channels.ClosedByInterruptException;
/**
* 使用BIO进行dump
*
* @author chuanyi
*/
public class BioSocketChannel implements SocketChannel {
static final int DEFAULT_CONNECT_TIMEOUT = 10 * 1000;
static final int SO_TIMEOUT = 1000;
private Socket socket;
private InputStream input;
private OutputStream output;
BioSocketChannel(Socket socket) throws IOException{
this.socket = socket;
this.input = socket.getInputStream();
this.output = socket.getOutputStream();
}
public void write(byte[]... buf) throws IOException {
OutputStream output = this.output;
if (output != null) {
for (byte[] bs : buf) {
output.write(bs);
}
} else {
throw new SocketException("Socket already closed.");
}
}
public byte[] read(int readSize) throws IOException {
InputStream input = this.input;
byte[] data = new byte[readSize];
int remain = readSize;
if (input == null) {
throw new SocketException("Socket already closed.");
}
while (remain > 0) {
try {
int read = input.read(data, readSize - remain, remain);
if (read > -1) {
remain -= read;
} else {
throw new IOException("EOF encountered.");
}
} catch (SocketTimeoutException te) {
if (Thread.interrupted()) {
throw new ClosedByInterruptException();
}
}
}
return data;
}
public byte[] read(int readSize, int timeout) throws IOException {
InputStream input = this.input;
byte[] data = new byte[readSize];
int remain = readSize;
int accTimeout = 0;
if (input == null) {
throw new SocketException("Socket already closed.");
}
while (remain > 0 && accTimeout < timeout) {
try {
int read = input.read(data, readSize - remain, remain);
if (read > -1) {
remain -= read;
} else {
throw new IOException("EOF encountered.");
}
} catch (SocketTimeoutException te) {
if (Thread.interrupted()) {
throw new ClosedByInterruptException();
}
accTimeout += SO_TIMEOUT;
}
}
if (remain > 0 && accTimeout >= timeout) {
throw new SocketTimeoutException("Timeout occurred, failed to read " + readSize + " bytes in " + timeout
+ " milliseconds.");
}
return data;
}
public boolean isConnected() {
Socket socket = this.socket;
if (socket != null) {
return socket.isConnected();
}
return false;
}
public SocketAddress getRemoteSocketAddress() {
Socket socket = this.socket;
if (socket != null) {
return socket.getRemoteSocketAddress();
}
return null;
}
public void close() {
Socket socket = this.socket;
if (socket != null) {
try {
socket.shutdownInput();
} catch (IOException e) {
// Ignore, could not do anymore
}
try {
socket.shutdownOutput();
} catch (IOException e) {
// Ignore, could not do anymore
}
try {
socket.close();
} catch (IOException e) {
// Ignore, could not do anymore
}
}
this.input = null;
this.output = null;
this.socket = null;
}
}
@@ -0,0 +1,24 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import java.net.Socket;
import java.net.SocketAddress;
/**
* @author luoyaogui 实现channel的管理(监听连接、读数据、回收) 2016-12-28
* @author chuanyi 2018-3-3 保留<code>open</code>减少文件变更数量
*/
public abstract class BioSocketChannelPool {
public static BioSocketChannel open(SocketAddress address) throws Exception {
Socket socket = new Socket();
socket.setReceiveBufferSize(32 * 1024);
socket.setSendBufferSize(32 * 1024);
socket.setSoTimeout(BioSocketChannel.SO_TIMEOUT);
socket.setTcpNoDelay(true);
socket.setKeepAlive(true);
socket.setReuseAddress(true);
socket.connect(address, BioSocketChannel.DEFAULT_CONNECT_TIMEOUT);
return new BioSocketChannel(socket);
}
}
@@ -0,0 +1,227 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.buffer.Unpooled;
import io.netty.channel.Channel;
import io.netty.util.internal.OutOfDirectMemoryError;
import io.netty.util.internal.PlatformDependent;
import io.netty.util.internal.SystemPropertyUtil;
import java.io.IOException;
import java.net.SocketAddress;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* 封装netty的通信channel和数据接收缓存,实现读、写、连接校验的功能。 2016-12-28
*
* @author luoyaogui
*/
public class NettySocketChannel implements SocketChannel {
private static final Logger logger = LoggerFactory.getLogger(SocketChannel.class);
private static final int WAIT_PERIOD = 10; // milliseconds
private static final int DEFAULT_INIT_BUFFER_SIZE = 1024 * 1024; // 1MB,默认初始缓存大小
// 参考 mysql-connector-java-5.1.40.jar: com.mysql.jdbc.MysqlIO.maxThreeBytes
// < 256 * 256 * 256 = 16MB
private static final int DEFAULT_MAX_BUFFER_SIZE = 16 * DEFAULT_INIT_BUFFER_SIZE; // 16MB,默认最大缓存大小
private Channel channel = null;
private Object lock = new Object();
private ByteBuf cache = PooledByteBufAllocator.DEFAULT.directBuffer(DEFAULT_INIT_BUFFER_SIZE); // 缓存大小
private int maxDirectBuffer = cache.maxCapacity();
public Channel getChannel() {
return channel;
}
public void setChannel(Channel channel) {
this.channel = channel;
}
public void writeCache(ByteBuf buf) throws InterruptedException, IOException {
synchronized (lock) {
while (true) {
if (null == cache) {
throw new IOException("socket is closed !");
}
// source buffer is empty.
if (!buf.isReadable()) {
break;
}
// 默认缓存大小不够用时需自动清理或扩充,否则将因缓存空间不足而造成I/O超时假象
int length = buf.readableBytes();
int deltaSize = length - cache.writableBytes();
if (deltaSize > 0) {
// 首先避免频繁分配内存(扩容/收缩),其次避免频繁移动内存(清理)
if (cache.readerIndex() >= deltaSize) { // 可以清理
// 回收已读空间,重置读写指针
cache.discardReadBytes();
// 恢复自动扩充的过大缓存到默认初始缓存大小,释放空间
int oldCapacity = cache.capacity();
if (oldCapacity > DEFAULT_MAX_BUFFER_SIZE) { // 尝试收缩
int newCapacity = cache.writerIndex();
newCapacity = ((newCapacity - 1) / DEFAULT_INIT_BUFFER_SIZE + 1) * DEFAULT_INIT_BUFFER_SIZE; // 对齐
int quarter = (newCapacity >> 2); // 至少留空四分之一
quarter = ((quarter - 1) / DEFAULT_INIT_BUFFER_SIZE + 1) * DEFAULT_INIT_BUFFER_SIZE; // 对齐
newCapacity += quarter; // 留空四分之一
if (newCapacity < (oldCapacity >> 1)) { // 至少收缩二分之一
try {
cache.capacity(newCapacity);
logger.info("shrink cache capacity: {} - {} = {} bytes",
oldCapacity,
oldCapacity - newCapacity,
newCapacity);
} catch (OutOfMemoryError ignore) {
maxDirectBuffer = oldCapacity; // 未来不再超过当前容量,记录日志后继续
logger.warn("cache OutOfMemoryError: {} bytes", newCapacity, ignore);
}
}
}
} else { // 尝试扩容
int oldCapacity = cache.capacity();
if (oldCapacity < maxDirectBuffer) {
int quarter = (oldCapacity >> 2); // 至少扩容四分之一
quarter = ((quarter - 1) / DEFAULT_INIT_BUFFER_SIZE + 1) * DEFAULT_INIT_BUFFER_SIZE; // 对齐
deltaSize = ((deltaSize - 1) / quarter + 1) * quarter; // 对齐
int newCapacity = oldCapacity + deltaSize;
if (newCapacity > maxDirectBuffer) {
newCapacity = maxDirectBuffer;
}
try {
cache.capacity(newCapacity);
logger.info("expand cache capacity: {} + {} = {} bytes",
oldCapacity,
newCapacity - oldCapacity,
newCapacity);
} catch (OutOfDirectMemoryError e) {
// failed to allocate 885571168 byte(s) of
// direct memory (used: 1002946176, max:
// 1888485376)
long maxDirectMemory = SystemPropertyUtil.getLong("io.netty.maxDirectMemory", -1);
if (maxDirectMemory < 0) {
maxDirectMemory = PlatformDependent.maxDirectMemory();
}
if (maxDirectBuffer > maxDirectMemory) {
maxDirectBuffer = (int) maxDirectMemory;
newCapacity = maxDirectBuffer;
logger.warn("resize maxDirectBuffer: {} bytes", maxDirectBuffer, e);
try {
cache.capacity(newCapacity);
logger.info("expand cache capacity: {} + {} = {} bytes",
oldCapacity,
newCapacity - oldCapacity,
newCapacity);
} catch (OutOfMemoryError ignore) {
maxDirectBuffer = oldCapacity; // 未来不再超过当前容量,记录日志后继续
logger.warn("cache OutOfMemoryError: {} bytes", newCapacity, ignore);
}
} else {
maxDirectBuffer = oldCapacity; // 未来不再超过当前容量,记录日志后继续
logger.warn("cache OutOfDirectMemoryError: {} bytes", newCapacity, e);
}
} catch (OutOfMemoryError ignore) {
maxDirectBuffer = oldCapacity; // 未来不再超过当前容量,记录日志后继续
logger.warn("cache OutOfMemoryError: {} bytes", newCapacity, ignore);
}
}
}
deltaSize = length - cache.writableBytes();
}
if (deltaSize != length) {
// deltaSize <= 0 可全部写入,deltaSize > 0 只能部分写入
if (deltaSize <= 0) {
cache.writeBytes(buf, length);
break;
} else {
cache.writeBytes(buf, length - deltaSize);
}
}
// dest buffer is full.
lock.wait(WAIT_PERIOD);
// 回收已读空间,重置读写指针
cache.discardReadBytes();
}
}
}
public void write(byte[]... buf) throws IOException {
if (channel != null && channel.isWritable()) {
channel.writeAndFlush(Unpooled.copiedBuffer(buf));
} else {
throw new IOException("write failed ! please checking !");
}
}
public byte[] read(int readSize) throws IOException {
return read(readSize, 0);
}
public byte[] read(int readSize, int timeout) throws IOException {
int accumulatedWaitTime = 0;
// 若读取内容较长,则自动扩充超时时间,以初始缓存大小为基准计算倍数
if (timeout > 0 && readSize > DEFAULT_INIT_BUFFER_SIZE) {
timeout *= (readSize / DEFAULT_INIT_BUFFER_SIZE + 1);
}
do {
if (readSize > cache.readableBytes()) {
if (null == channel) {
throw new IOException("socket has Interrupted !");
}
if (timeout > 0) {
accumulatedWaitTime += WAIT_PERIOD;
if (accumulatedWaitTime > timeout) {
StringBuilder sb = new StringBuilder("socket read timeout occured !");
sb.append(" readSize = ").append(readSize);
sb.append(", readableBytes = ").append(cache.readableBytes());
sb.append(", timeout = ").append(timeout);
throw new IOException(sb.toString());
}
}
synchronized (this) {
try {
wait(WAIT_PERIOD);
} catch (InterruptedException e) {
throw new IOException("socket has Interrupted !");
}
}
} else {
byte[] back = new byte[readSize];
synchronized (lock) {
cache.readBytes(back);
}
return back;
}
} while (true);
}
public boolean isConnected() {
return channel != null ? true : false;
}
public SocketAddress getRemoteSocketAddress() {
return channel != null ? channel.remoteAddress() : null;
}
public void close() {
if (channel != null) {
channel.close();
}
channel = null;
// A fatal error has been detected by the Java Runtime Environment:
// EXCEPTION_ACCESS_VIOLATION (0xc0000005)
synchronized (lock) {
cache.discardReadBytes();// 回收已占用的内存
cache.release();// 释放整个内存
cache = null;
}
}
}
@@ -0,0 +1,111 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.AdaptiveRecvByteBufAllocator;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelOption;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioSocketChannel;
import java.io.IOException;
import java.net.SocketAddress;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author luoyaogui 实现channel的管理(监听连接、读数据、回收) 2016-12-28
*/
@SuppressWarnings({ "rawtypes", "deprecation" })
public abstract class NettySocketChannelPool {
private static EventLoopGroup group = new NioEventLoopGroup(); // 非阻塞IO线程组
private static Bootstrap boot = new Bootstrap(); // 主
private static Map<Channel, SocketChannel> chManager = new ConcurrentHashMap<Channel, SocketChannel>();
private static final Logger logger = LoggerFactory.getLogger(NettySocketChannelPool.class);
static {
boot.group(group)
.channel(NioSocketChannel.class)
.option(ChannelOption.SO_RCVBUF, 32 * 1024)
.option(ChannelOption.SO_SNDBUF, 32 * 1024)
.option(ChannelOption.TCP_NODELAY, true)
// 如果是延时敏感型应用,建议关闭Nagle算法
.option(ChannelOption.SO_KEEPALIVE, true)
.option(ChannelOption.RCVBUF_ALLOCATOR, AdaptiveRecvByteBufAllocator.DEFAULT)
.option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
//
.handler(new ChannelInitializer() {
@Override
protected void initChannel(Channel ch) throws Exception {
ch.pipeline().addLast(new BusinessHandler());// 命令过滤和handler添加管理
}
});
}
public static SocketChannel open(SocketAddress address) throws Exception {
SocketChannel socket = null;
ChannelFuture future = boot.connect(address).sync();
if (future.isSuccess()) {
future.channel().pipeline().get(BusinessHandler.class).latch.await();
socket = chManager.get(future.channel());
}
if (null == socket) {
throw new IOException("can't create socket!");
}
return socket;
}
public static class BusinessHandler extends SimpleChannelInboundHandler<ByteBuf> {
private NettySocketChannel socket = null;
private final CountDownLatch latch = new CountDownLatch(1);
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
socket.setChannel(null);
chManager.remove(ctx.channel());// 移除
super.channelInactive(ctx);
}
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
socket = new NettySocketChannel();
socket.setChannel(ctx.channel());
chManager.put(ctx.channel(), socket);
latch.countDown();
super.channelActive(ctx);
}
@Override
protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {
if (socket != null) {
socket.writeCache(msg);
} else {
// TODO: need graceful error handler.
logger.error("no socket available.");
}
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
// need output error for troubeshooting.
logger.error("business error.", cause);
ctx.close();
}
}
}
@@ -1,85 +1,23 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.buffer.Unpooled;
import io.netty.channel.Channel;
import java.io.IOException;
import java.net.SocketAddress;
/**
* 封装netty的通信channel和数据接收缓存,实现读、写、连接校验的功能。 2016-12-28
*
* @author luoyaogui
* @author agapple 2018年3月12日 下午10:36:44
* @since 1.0.26
*/
public class SocketChannel {
public interface SocketChannel {
private Channel channel = null;
private Object lock = new Object();
private ByteBuf cache = PooledByteBufAllocator.DEFAULT.directBuffer(1024 * 1024); // 缓存大小
public void write(byte[]... buf) throws IOException;
public Channel getChannel() {
return channel;
}
public byte[] read(int readSize) throws IOException;
public void setChannel(Channel channel) {
this.channel = channel;
}
public byte[] read(int readSize, int timeout) throws IOException;
public void writeCache(ByteBuf buf) {
synchronized (lock) {
cache.discardReadBytes();// 回收内存
cache.writeBytes(buf);
}
}
public boolean isConnected();
public void writeChannel(byte[]... buf) throws IOException {
if (channel != null && channel.isWritable()) {
channel.writeAndFlush(Unpooled.copiedBuffer(buf));
} else {
throw new IOException("write failed ! please checking !");
}
}
public SocketAddress getRemoteSocketAddress();
public byte[] read(int readSize) throws IOException {
do {
if (readSize > cache.readableBytes()) {
if (null == channel) {
throw new java.nio.channels.ClosedByInterruptException();
}
synchronized (this) {
try {
wait(100);
} catch (InterruptedException e) {
throw new java.nio.channels.ClosedByInterruptException();
}
}
} else {
byte[] back = new byte[readSize];
synchronized (lock) {
cache.readBytes(back);
}
return back;
}
} while (true);
}
public boolean isConnected() {
return channel != null ? true : false;
}
public SocketAddress getRemoteSocketAddress() {
return channel != null ? channel.remoteAddress() : null;
}
public void close() {
if (channel != null) {
channel.close();
}
channel = null;
cache.discardReadBytes();// 回收已占用的内存
cache.release();// 释放整个内存
cache = null;
}
public void close();
}
@@ -1,103 +1,35 @@
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.AdaptiveRecvByteBufAllocator;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelOption;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioSocketChannel;
import io.netty.util.ReferenceCountUtil;
import java.io.IOException;
import java.net.SocketAddress;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import com.alibaba.otter.canal.common.utils.BooleanMutex;
/**
* @author luoyaogui 实现channel的管理(监听连接、读数据、回收) 2016-12-28
*/
@SuppressWarnings({ "rawtypes", "deprecation" })
public abstract class SocketChannelPool {
private static EventLoopGroup group = new NioEventLoopGroup(); // 非阻塞IO线程组
private static Bootstrap boot = new Bootstrap(); // 主
private static Map<Channel, SocketChannel> chManager = new ConcurrentHashMap<Channel, SocketChannel>();
static {
boot.group(group)
.channel(NioSocketChannel.class)
.option(ChannelOption.SO_RCVBUF, 32 * 1024)
.option(ChannelOption.SO_SNDBUF, 32 * 1024)
.option(ChannelOption.TCP_NODELAY, true)
// 如果是延时敏感型应用,建议关闭Nagle算法
.option(ChannelOption.SO_KEEPALIVE, true)
.option(ChannelOption.RCVBUF_ALLOCATOR, AdaptiveRecvByteBufAllocator.DEFAULT)
.option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
//
.handler(new ChannelInitializer() {
@Override
protected void initChannel(Channel arg0) throws Exception {
arg0.pipeline().addLast(new BusinessHandler());// 命令过滤和handler添加管理
}
});
}
public static SocketChannel open(SocketAddress address) throws Exception {
final SocketChannel socket = new SocketChannel();
final BooleanMutex mutex = new BooleanMutex(false);
boot.connect(address).addListener(new ChannelFutureListener() {
@Override
public void operationComplete(ChannelFuture arg0) throws Exception {
if (arg0.isSuccess()) {
socket.setChannel(arg0.channel());
}
mutex.set(true);
}
});
// wait for complete
mutex.get();
if (null == socket.getChannel()) {
throw new IOException("can't create socket!");
}
chManager.put(socket.getChannel(), socket);
return socket;
}
public static class BusinessHandler extends ChannelInboundHandlerAdapter {
private SocketChannel socket = null;
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
socket.setChannel(null);
chManager.remove(ctx.channel());// 移除
}
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
if (null == socket) socket = chManager.get(ctx.channel());
if (socket != null) {
socket.writeCache((ByteBuf) msg);
}
ReferenceCountUtil.release(msg);// 添加防止内存泄漏的
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
ctx.close();
}
}
}
package com.alibaba.otter.canal.parse.driver.mysql.socket;
import java.net.SocketAddress;
import org.apache.commons.lang.StringUtils;
/**
* @author agapple 2018年3月12日 下午10:46:22
* @since 1.0.26
*/
public abstract class SocketChannelPool {
public static SocketChannel open(SocketAddress address) throws Exception {
String type = chooseSocketChannel();
if ("netty".equalsIgnoreCase(type)) {
return NettySocketChannelPool.open(address);
} else {
return BioSocketChannelPool.open(address);
}
}
private static String chooseSocketChannel() {
String socketChannel = System.getenv("canal.socketChannel");
if (StringUtils.isEmpty(socketChannel)) {
socketChannel = System.getProperty("canal.socketChannel");
}
if (StringUtils.isEmpty(socketChannel)) {
socketChannel = "bio"; // bio or netty
}
return socketChannel;
}
}
@@ -108,7 +108,18 @@ public abstract class ByteHelper {
return out.toByteArray();
}
public static void write8ByteUnsignedIntLittleEndian(long data, ByteArrayOutputStream out) {
out.write((byte) (data & 0xFF));
out.write((byte) (data >>> 8));
out.write((byte) (data >>> 16));
out.write((byte) (data >>> 24));
out.write((byte) (data >>> 32));
out.write((byte) (data >>> 40));
out.write((byte) (data >>> 48));
out.write((byte) (data >>> 56));
}
public static void writeUnsignedIntLittleEndian(long data, ByteArrayOutputStream out) {
out.write((byte) (data & 0xFF));
out.write((byte) (data >>> 8));
@@ -13,12 +13,22 @@ public abstract class PacketManager {
return header;
}
public static HeaderPacket readHeader(SocketChannel ch, int len, int timeout) throws IOException {
HeaderPacket header = new HeaderPacket();
header.fromBytes(ch.read(len, timeout));
return header;
}
public static byte[] readBytes(SocketChannel ch, int len) throws IOException {
return ch.read(len);
}
public static byte[] readBytes(SocketChannel ch, int len, int timeout) throws IOException {
return ch.read(len, timeout);
}
public static void writePkg(SocketChannel ch, byte[]... srcs) throws IOException {
ch.writeChannel(srcs);
ch.write(srcs);
}
public static void writeBody(SocketChannel ch, byte[] body) throws IOException {
@@ -29,6 +39,6 @@ public abstract class PacketManager {
HeaderPacket header = new HeaderPacket();
header.setPacketBodyLength(body.length);
header.setPacketSequenceNumber(packetSeqNumber);
ch.writeChannel(header.toBytes(), body);
ch.write(header.toBytes(), body);
}
}
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -56,10 +56,10 @@ public class AbstractCanalClientTest {
context_format += "****************************************************" + SEP;
row_format = SEP
+ "----------------> binlog[{}:{}] , name[{},{}] , eventType : {} , executeTime : {} , delay : {}ms"
+ "----------------> binlog[{}:{}] , name[{},{}] , eventType : {} , executeTime : {}({}) , delay : {} ms"
+ SEP;
transaction_format = SEP + "================> binlog[{}:{}] , executeTime : {} , delay : {}ms" + SEP;
transaction_format = SEP + "================> binlog[{}:{}] , executeTime : {}({}) , delay : {}ms" + SEP;
}
@@ -165,6 +165,8 @@ public class AbstractCanalClientTest {
for (Entry entry : entrys) {
long executeTime = entry.getHeader().getExecuteTime();
long delayTime = new Date().getTime() - executeTime;
Date date = new Date(entry.getHeader().getExecuteTime());
SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
if (entry.getEntryType() == EntryType.TRANSACTIONBEGIN || entry.getEntryType() == EntryType.TRANSACTIONEND) {
if (entry.getEntryType() == EntryType.TRANSACTIONBEGIN) {
@@ -178,7 +180,8 @@ public class AbstractCanalClientTest {
logger.info(transaction_format,
new Object[] { entry.getHeader().getLogfileName(),
String.valueOf(entry.getHeader().getLogfileOffset()),
String.valueOf(entry.getHeader().getExecuteTime()), String.valueOf(delayTime) });
String.valueOf(entry.getHeader().getExecuteTime()), simpleDateFormat.format(date),
String.valueOf(delayTime) });
logger.info(" BEGIN ----> Thread id: {}", begin.getThreadId());
} else if (entry.getEntryType() == EntryType.TRANSACTIONEND) {
TransactionEnd end = null;
@@ -193,7 +196,8 @@ public class AbstractCanalClientTest {
logger.info(transaction_format,
new Object[] { entry.getHeader().getLogfileName(),
String.valueOf(entry.getHeader().getLogfileOffset()),
String.valueOf(entry.getHeader().getExecuteTime()), String.valueOf(delayTime) });
String.valueOf(entry.getHeader().getExecuteTime()), simpleDateFormat.format(date),
String.valueOf(delayTime) });
}
continue;
@@ -213,7 +217,8 @@ public class AbstractCanalClientTest {
new Object[] { entry.getHeader().getLogfileName(),
String.valueOf(entry.getHeader().getLogfileOffset()), entry.getHeader().getSchemaName(),
entry.getHeader().getTableName(), eventType,
String.valueOf(entry.getHeader().getExecuteTime()), String.valueOf(delayTime) });
String.valueOf(entry.getHeader().getExecuteTime()), simpleDateFormat.format(date),
String.valueOf(delayTime) });
if (eventType == EventType.QUERY || rowChage.getIsDdl()) {
logger.info(" sql ----> " + rowChage.getSql() + SEP);
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../../pom.xml</relativePath>
</parent>
<artifactId>canal.instance.core</artifactId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+6 -3
View File
@@ -1,10 +1,9 @@
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<artifactId>canal.parse</artifactId>
@@ -50,6 +49,10 @@
<groupId>com.alibaba</groupId>
<artifactId>druid</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba.fastsql</groupId>
<artifactId>fastsql</artifactId>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
@@ -171,7 +171,7 @@ public abstract class AbstractEventParser<EVENT> extends AbstractCanalLifeCycle
throw new CanalParseException("can't find init table meta for " + destination
+ " with position : " + startPosition);
}
logger.info("find start position : {}", startPosition.toString());
logger.warn("find start position : {}", startPosition.toString());
// 重新链接,因为在找position过程中可能有状态,需要断开后重建
erosaConnection.reconnect();
@@ -3,9 +3,10 @@ package com.alibaba.otter.canal.parse.inbound;
import java.util.ArrayList;
import java.util.List;
import com.taobao.tddl.dbsync.binlog.event.TableMapLogEvent;
import org.apache.commons.lang.StringUtils;
import com.taobao.tddl.dbsync.binlog.event.TableMapLogEvent;
/**
* 描述数据meta对象,mysql binlog中对应的{@linkplain TableMapLogEvent}包含的信息不全
*
@@ -127,6 +128,7 @@ public class TableMeta {
private boolean key;
private String defaultValue;
private String extra;
private boolean unique;
public String getColumnName() {
return columnName;
@@ -180,10 +182,21 @@ public class TableMeta {
this.extra = extra;
}
public String toString() {
return "FieldMeta [columnName=" + columnName + ", columnType=" + columnType + ", defaultValue="
+ defaultValue + ", nullable=" + nullable + ", key=" + key + "]";
public boolean isUnique() {
return unique;
}
public void setUnique(boolean unique) {
this.unique = unique;
}
@Override
public String toString() {
return "FieldMeta [columnName=" + columnName + ", columnType=" + columnType + ", nullable=" + nullable
+ ", key=" + key + ", defaultValue=" + defaultValue + ", extra=" + extra + ", unique=" + unique
+ "]";
}
}
}
@@ -159,10 +159,11 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
public void setTsdbSpringXml(String tsdbSpringXml) {
this.tsdbSpringXml = tsdbSpringXml;
if (tableMetaTSDB == null) {
// 初始化
tableMetaTSDB = TableMetaTSDBBuilder.build(destination, tsdbSpringXml);
if (this.enableTsdb) {
if (tableMetaTSDB == null) {
// 初始化
tableMetaTSDB = TableMetaTSDBBuilder.build(destination, tsdbSpringXml);
}
}
}
@@ -1,9 +1,12 @@
package com.alibaba.otter.canal.parse.inbound.mysql;
import static com.alibaba.otter.canal.parse.inbound.mysql.dbsync.DirectLogFetcher.MASTER_HEARTBEAT_PERIOD_SECONDS;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.charset.Charset;
import java.util.List;
import java.util.concurrent.TimeUnit;
import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
@@ -14,6 +17,9 @@ import com.alibaba.otter.canal.parse.driver.mysql.MysqlQueryExecutor;
import com.alibaba.otter.canal.parse.driver.mysql.MysqlUpdateExecutor;
import com.alibaba.otter.canal.parse.driver.mysql.packets.HeaderPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.client.BinlogDumpCommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.client.RegisterSlaveCommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.client.SemiAckCommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.server.ErrorPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.server.ResultSetPacket;
import com.alibaba.otter.canal.parse.driver.mysql.utils.PacketManager;
import com.alibaba.otter.canal.parse.exception.CanalParseException;
@@ -129,6 +135,7 @@ public class MysqlConnection implements ErosaConnection {
public void dump(String binlogfilename, Long binlogPosition, SinkFunction func) throws IOException {
updateSettings();
sendRegisterSlave();
sendBinlogDump(binlogfilename, binlogPosition);
DirectLogFetcher fetcher = new DirectLogFetcher(connector.getReceiveBufferSize());
fetcher.start(connector.getChannel());
@@ -145,6 +152,10 @@ public class MysqlConnection implements ErosaConnection {
if (!func.sink(event)) {
break;
}
if (event.getSemival() == 1) {
sendSemiAck(context.getLogPosition().getFileName(), binlogPosition);
}
}
}
@@ -152,11 +163,40 @@ public class MysqlConnection implements ErosaConnection {
throw new NullPointerException("Not implement yet");
}
private void sendRegisterSlave() throws IOException {
RegisterSlaveCommandPacket cmd = new RegisterSlaveCommandPacket();
cmd.reportHost = authInfo.getAddress().getAddress().getHostAddress();
cmd.reportPasswd = authInfo.getPassword();
cmd.reportUser = authInfo.getUsername();
cmd.reportPort = authInfo.getAddress().getPort(); // 暂时先用master节点的port
cmd.serverId = this.slaveId;
byte[] cmdBody = cmd.toBytes();
logger.info("Register slave {}", cmd);
HeaderPacket header = new HeaderPacket();
header.setPacketBodyLength(cmdBody.length);
header.setPacketSequenceNumber((byte) 0x00);
PacketManager.writePkg(connector.getChannel(), header.toBytes(), cmdBody);
header = PacketManager.readHeader(connector.getChannel(), 4);
byte[] body = PacketManager.readBytes(connector.getChannel(), header.getPacketBodyLength());
assert body != null;
if (body[0] < 0) {
if (body[0] == -1) {
ErrorPacket err = new ErrorPacket();
err.fromBytes(body);
throw new IOException("Error When doing Register slave:" + err.toString());
} else {
throw new IOException("unpexpected packet with field_count=" + body[0]);
}
}
}
private void sendBinlogDump(String binlogfilename, Long binlogPosition) throws IOException {
BinlogDumpCommandPacket binlogDumpCmd = new BinlogDumpCommandPacket();
binlogDumpCmd.binlogFileName = binlogfilename;
binlogDumpCmd.binlogPosition = binlogPosition;
// binlogDumpCmd.slaveServerId = this.slaveId;
binlogDumpCmd.slaveServerId = this.slaveId;
byte[] cmdBody = binlogDumpCmd.toBytes();
@@ -168,6 +208,20 @@ public class MysqlConnection implements ErosaConnection {
connector.setDumping(true);
}
private void sendSemiAck(String binlogfilename, Long binlogPosition) throws IOException {
SemiAckCommandPacket semiAckCmd = new SemiAckCommandPacket();
semiAckCmd.binlogFileName = binlogfilename;
semiAckCmd.binlogPosition = binlogPosition;
byte[] cmdBody = semiAckCmd.toBytes();
logger.info("SEMI ACK with position:{}", semiAckCmd);
HeaderPacket semiAckHeader = new HeaderPacket();
semiAckHeader.setPacketBodyLength(cmdBody.length);
semiAckHeader.setPacketSequenceNumber((byte) 0x00);
PacketManager.writePkg(connector.getChannel(), semiAckHeader.toBytes(), cmdBody);
}
public MysqlConnection fork() {
MysqlConnection connection = new MysqlConnection();
connection.setCharset(getCharset());
@@ -242,6 +296,22 @@ public class MysqlConnection implements ErosaConnection {
} catch (Exception e) {
logger.warn("update mariadb_slave_capability failed", e);
}
/**
* MASTER_HEARTBEAT_PERIOD sets the interval in seconds between
* replication heartbeats. Whenever the master's binary log is updated
* with an event, the waiting period for the next heartbeat is reset.
* interval is a decimal value having the range 0 to 4294967 seconds and
* a resolution in milliseconds; the smallest nonzero value is 0.001.
* Heartbeats are sent by the master only if there are no unsent events
* in the binary log file for a period longer than interval.
*/
try {
long periodNano = TimeUnit.SECONDS.toNanos(MASTER_HEARTBEAT_PERIOD_SECONDS);
update("SET @master_heartbeat_period=" + periodNano);
} catch (Exception e) {
logger.warn("update master_heartbeat_period failed", e);
}
}
/**
@@ -8,7 +8,6 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.TimerTask;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import org.apache.commons.lang.StringUtils;
@@ -374,7 +373,8 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
return findAsPerTimestampInSpecificLogFile(mysqlConnection,
startTimestamp,
endPosition,
endPosition.getJournalName());
endPosition.getJournalName(),
true);
} else {
return endPosition;
}
@@ -385,10 +385,16 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
if (tableMetaTSDB != null && (fixedPosition.getTimestamp() == null || fixedPosition.getTimestamp() <= 0)) {
// 使用一个未来极大的时间,基于位点进行定位
long startTimestamp = System.currentTimeMillis() + 102L * 365 * 24 * 3600 * 1000; // 当前时间的未来102年
return findAsPerTimestampInSpecificLogFile(mysqlConnection,
EntryPosition entryPosition = findAsPerTimestampInSpecificLogFile(mysqlConnection,
startTimestamp,
fixedPosition,
fixedPosition.getJournalName());
fixedPosition.getJournalName(),
true);
if (entryPosition == null) {
throw new CanalParseException("[fixed timestamp] can't found begin/commit position before with fixed position"
+ fixedPosition.getJournalName() + ":" + fixedPosition.getPosition());
}
return entryPosition;
} else {
return fixedPosition;
}
@@ -441,7 +447,8 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
specificLogFilePosition = findAsPerTimestampInSpecificLogFile(mysqlConnection,
entryPosition.getTimestamp(),
endPosition,
entryPosition.getJournalName());
entryPosition.getJournalName(),
true);
}
}
@@ -492,10 +499,10 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
// 主要考虑一个事务执行时间可能会几秒种,如果仅仅按照timestamp相同,则可能会丢失事务的前半部分数据
private Long findTransactionBeginPosition(ErosaConnection mysqlConnection, final EntryPosition entryPosition)
throws IOException {
// 尝试找到一个合适的位置
final AtomicBoolean reDump = new AtomicBoolean(false);
// 针对开始的第一条为非Begin记录,需要从该binlog扫描
final AtomicLong preTransactionStartPosition = new AtomicLong(0L);
mysqlConnection.reconnect();
mysqlConnection.seek(entryPosition.getJournalName(), entryPosition.getPosition(), new SinkFunction<LogEvent>() {
mysqlConnection.seek(entryPosition.getJournalName(), 4L, new SinkFunction<LogEvent>() {
private LogPosition lastPosition;
@@ -506,69 +513,33 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
return true;
}
// 直接查询第一条业务数据,确认是否为事务Begin/End
if (CanalEntry.EntryType.TRANSACTIONBEGIN == entry.getEntryType()
|| CanalEntry.EntryType.TRANSACTIONEND == entry.getEntryType()) {
lastPosition = buildLastPosition(entry);
return false;
} else {
reDump.set(true);
lastPosition = buildLastPosition(entry);
return false;
// 直接查询第一条业务数据,确认是否为事务Begin
// 记录一下transaction begin position
if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN
&& entry.getHeader().getLogfileOffset() < entryPosition.getPosition()) {
preTransactionStartPosition.set(entry.getHeader().getLogfileOffset());
}
if (entry.getHeader().getLogfileOffset() >= entryPosition.getPosition()) {
return false;// 退出
}
lastPosition = buildLastPosition(entry);
} catch (Exception e) {
// 上一次记录的poistion可能为一条update/insert/delete变更事件,直接进行dump的话,会缺少tableMap事件,导致tableId未进行解析
processSinkError(e, lastPosition, entryPosition.getJournalName(), entryPosition.getPosition());
reDump.set(true);
return false;
}
return running;
}
});
// 针对开始的第一条为非Begin记录,需要从该binlog扫描
if (reDump.get()) {
final AtomicLong preTransactionStartPosition = new AtomicLong(0L);
mysqlConnection.reconnect();
mysqlConnection.seek(entryPosition.getJournalName(), 4L, new SinkFunction<LogEvent>() {
private LogPosition lastPosition;
public boolean sink(LogEvent event) {
try {
CanalEntry.Entry entry = parseAndProfilingIfNecessary(event, true);
if (entry == null) {
return true;
}
// 直接查询第一条业务数据,确认是否为事务Begin
// 记录一下transaction begin position
if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN
&& entry.getHeader().getLogfileOffset() < entryPosition.getPosition()) {
preTransactionStartPosition.set(entry.getHeader().getLogfileOffset());
}
if (entry.getHeader().getLogfileOffset() >= entryPosition.getPosition()) {
return false;// 退出
}
lastPosition = buildLastPosition(entry);
} catch (Exception e) {
processSinkError(e, lastPosition, entryPosition.getJournalName(), entryPosition.getPosition());
return false;
}
return running;
}
});
// 判断一下找到的最接近position的事务头的位置
if (preTransactionStartPosition.get() > entryPosition.getPosition()) {
logger.error("preTransactionEndPosition greater than startPosition from zk or localconf, maybe lost data");
throw new CanalParseException("preTransactionStartPosition greater than startPosition from zk or localconf, maybe lost data");
}
return preTransactionStartPosition.get();
} else {
return entryPosition.getPosition();
// 判断一下找到的最接近position的事务头的位置
if (preTransactionStartPosition.get() > entryPosition.getPosition()) {
logger.error("preTransactionEndPosition greater than startPosition from zk or localconf, maybe lost data");
throw new CanalParseException("preTransactionStartPosition greater than startPosition from zk or localconf, maybe lost data");
}
return preTransactionStartPosition.get();
}
// 根据时间查找binlog位置
@@ -585,7 +556,8 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
EntryPosition entryPosition = findAsPerTimestampInSpecificLogFile(mysqlConnection,
startTimestamp,
endPosition,
startSearchBinlogFile);
startSearchBinlogFile,
false);
if (entryPosition == null) {
if (StringUtils.equalsIgnoreCase(minBinlogFileName, startSearchBinlogFile)) {
// 已经找到最早的一个binlog,没必要往前找了
@@ -733,7 +705,8 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
private EntryPosition findAsPerTimestampInSpecificLogFile(MysqlConnection mysqlConnection,
final Long startTimestamp,
final EntryPosition endPosition,
final String searchBinlogFile) {
final String searchBinlogFile,
final Boolean justForPositionTimestamp) {
final LogPosition logPosition = new LogPosition();
try {
@@ -747,6 +720,15 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
EntryPosition entryPosition = null;
try {
CanalEntry.Entry entry = parseAndProfilingIfNecessary(event, true);
if (justForPositionTimestamp && logPosition.getPostion() == null && event.getWhen() > 0) {
// 初始位点
entryPosition = new EntryPosition(searchBinlogFile,
event.getLogPos(),
event.getWhen() * 1000,
event.getServerId());
logPosition.setPostion(entryPosition);
}
if (entry == null) {
return true;
}
@@ -758,8 +740,10 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
if (CanalEntry.EntryType.TRANSACTIONBEGIN.equals(entry.getEntryType())
|| CanalEntry.EntryType.TRANSACTIONEND.equals(entry.getEntryType())) {
logger.debug("compare exit condition:{},{},{}, startTimestamp={}...", new Object[] {
logfilename, logfileoffset, logposTimestamp, startTimestamp });
if (logger.isDebugEnabled()) {
logger.debug("compare exit condition:{},{},{}, startTimestamp={}...", new Object[] {
logfilename, logfileoffset, logposTimestamp, startTimestamp });
}
// 事务头和尾寻找第一条记录时间戳,如果最小的一条记录都不满足条件,可直接退出
if (logposTimestamp >= startTimestamp) {
return false;
@@ -767,7 +751,7 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
}
if (StringUtils.equals(endPosition.getJournalName(), logfilename)
&& endPosition.getPosition() <= logfileoffset) {
&& endPosition.getPosition() < logfileoffset) {
return false;
}
@@ -776,14 +760,18 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
// data.length,代表该事务的下一条offest,避免多余的事务重复
if (CanalEntry.EntryType.TRANSACTIONEND.equals(entry.getEntryType())) {
entryPosition = new EntryPosition(logfilename, logfileoffset, logposTimestamp, serverId);
logger.debug("set {} to be pending start position before finding another proper one...",
entryPosition);
if (logger.isDebugEnabled()) {
logger.debug("set {} to be pending start position before finding another proper one...",
entryPosition);
}
logPosition.setPostion(entryPosition);
} else if (CanalEntry.EntryType.TRANSACTIONBEGIN.equals(entry.getEntryType())) {
// 当前事务开始位点
entryPosition = new EntryPosition(logfilename, logfileoffset, logposTimestamp, serverId);
logger.debug("set {} to be pending start position before finding another proper one...",
entryPosition);
if (logger.isDebugEnabled()) {
logger.debug("set {} to be pending start position before finding another proper one...",
entryPosition);
}
logPosition.setPostion(entryPosition);
}
@@ -21,6 +21,11 @@ public class DirectLogFetcher extends LogFetcher {
protected static final Logger logger = LoggerFactory.getLogger(DirectLogFetcher.class);
// Master heartbeat interval
public static final int MASTER_HEARTBEAT_PERIOD_SECONDS = 15;
// +10s 确保 timeout > heartbeat interval
private static final int READ_TIMEOUT_MILLISECONDS = (MASTER_HEARTBEAT_PERIOD_SECONDS + 10) * 1000;
/** Command to dump binlog */
public static final byte COM_BINLOG_DUMP = 18;
@@ -37,6 +42,8 @@ public class DirectLogFetcher extends LogFetcher {
private SocketChannel channel;
private boolean issemi = false;
// private BufferedInputStream input;
public DirectLogFetcher(){
@@ -53,6 +60,10 @@ public class DirectLogFetcher extends LogFetcher {
public void start(SocketChannel channel) throws IOException {
this.channel = channel;
String dbsemi = System.getProperty("db.semi");
if ("1".equals(dbsemi)) {
issemi = true;
}
// 和mysql driver一样,提供buffer机制,提升读取binlog速度
// this.input = new
// BufferedInputStream(channel.socket().getInputStream(), 16384);
@@ -106,6 +117,14 @@ public class DirectLogFetcher extends LogFetcher {
}
}
// if mysql is in semi mode
if (issemi) {
// parse semi mark
int semimark = getUint8(NET_HEADER_SIZE + 1);
int semival = getUint8(NET_HEADER_SIZE + 2);
this.semival = semival;
}
// The first packet is a multi-packet, concatenate the packets.
while (netlen == MAX_PACKET_LENGTH) {
if (!fetch0(0, NET_HEADER_SIZE)) {
@@ -122,7 +141,11 @@ public class DirectLogFetcher extends LogFetcher {
}
// Preparing buffer variables to decoding.
origin = NET_HEADER_SIZE + 1;
if (issemi) {
origin = NET_HEADER_SIZE + 3;
} else {
origin = NET_HEADER_SIZE + 1;
}
position = origin;
limit -= origin;
return true;
@@ -148,7 +171,7 @@ public class DirectLogFetcher extends LogFetcher {
private final boolean fetch0(final int off, final int len) throws IOException {
ensureCapacity(off + len);
byte[] read = channel.read(len);
byte[] read = channel.read(len, READ_TIMEOUT_MILLISECONDS);
System.arraycopy(read, 0, this.buffer, off, len);
if (limit < off + len) limit = off + len;
@@ -364,6 +364,7 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
throw new TableIdNotFoundException("not found tableId:" + event.getTableId());
}
boolean isHeartBeat = isAliSQLHeartBeat(table.getDbName(), table.getTableName());
boolean isRDSHeartBeat = tableMetaCache.isOnRDS()
&& isRDSHeartBeat(table.getDbName(), table.getTableName());
@@ -387,6 +388,12 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
FieldMeta idMeta = new FieldMeta("id", "bigint(20)", true, false, "0");
FieldMeta typeMeta = new FieldMeta("type", "char(1)", false, true, "0");
tableMeta = new TableMeta(table.getDbName(), table.getTableName(), Arrays.asList(idMeta, typeMeta));
} else if (isHeartBeat) {
// 处理alisql模式的test.heartbeat心跳数据
// 心跳表基本无权限,需要mock一个tableMeta
FieldMeta idMeta = new FieldMeta("id", "smallint(6)", false, true, null);
FieldMeta typeMeta = new FieldMeta("type", "int(11)", true, false, null);
tableMeta = new TableMeta(table.getDbName(), table.getTableName(), Arrays.asList(idMeta, typeMeta));
}
EventType eventType = null;
@@ -766,6 +773,10 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
return "LONGTEXT".equalsIgnoreCase(columnType) || "MEDIUMTEXT".equalsIgnoreCase(columnType)
|| "TEXT".equalsIgnoreCase(columnType) || "TINYTEXT".equalsIgnoreCase(columnType);
}
private boolean isAliSQLHeartBeat(String schema, String table) {
return "test".equalsIgnoreCase(schema) && "heartbeat".equalsIgnoreCase(table);
}
private boolean isRDSHeartBeat(String schema, String table) {
return "mysql".equalsIgnoreCase(schema) && "ha_health_check".equalsIgnoreCase(table);
@@ -14,6 +14,9 @@ import com.alibaba.otter.canal.parse.exception.CanalParseException;
import com.alibaba.otter.canal.parse.inbound.TableMeta;
import com.alibaba.otter.canal.parse.inbound.TableMeta.FieldMeta;
import com.alibaba.otter.canal.parse.inbound.mysql.MysqlConnection;
import com.alibaba.otter.canal.parse.inbound.mysql.ddl.DruidDdlParser;
import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.DatabaseTableMeta;
import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.MemoryTableMeta;
import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.TableMetaTSDB;
import com.alibaba.otter.canal.protocol.position.EntryPosition;
import com.google.common.cache.CacheBuilder;
@@ -52,7 +55,7 @@ public class TableMetaCache {
public TableMeta load(String name) throws Exception {
try {
return getTableMetaByDB(name);
} catch (CanalParseException e) {
} catch (Throwable e) {
// 尝试做一次retry操作
try {
connection.reconnect();
@@ -76,16 +79,38 @@ public class TableMetaCache {
}
private TableMeta getTableMetaByDB(String fullname) throws IOException {
ResultSetPacket packet = connection.query("desc " + fullname);
String[] names = StringUtils.split(fullname, "`.`");
String schema = names[0];
String table = names[1].substring(0, names[1].length());
return new TableMeta(schema, table, parserTableMeta(packet));
try {
ResultSetPacket packet = connection.query("show create table " + fullname);
String[] names = StringUtils.split(fullname, "`.`");
String schema = names[0];
String table = names[1].substring(0, names[1].length());
return new TableMeta(schema, table, parseTableMeta(schema, table, packet));
} catch (Throwable e) { // fallback to desc table
ResultSetPacket packet = connection.query("desc " + fullname);
String[] names = StringUtils.split(fullname, "`.`");
String schema = names[0];
String table = names[1].substring(0, names[1].length());
return new TableMeta(schema, table, parseTableMetaByDesc(packet));
}
}
public static List<FieldMeta> parserTableMeta(ResultSetPacket packet) {
Map<String, Integer> nameMaps = new HashMap<String, Integer>(6, 1f);
public static List<FieldMeta> parseTableMeta(String schema, String table, ResultSetPacket packet) {
if (packet.getFieldValues().size() > 1) {
String createDDL = packet.getFieldValues().get(1);
MemoryTableMeta memoryTableMeta = new MemoryTableMeta();
memoryTableMeta.apply(DatabaseTableMeta.INIT_POSITION, schema, createDDL, null);
TableMeta tableMeta = memoryTableMeta.find(schema, table);
return tableMeta.getFields();
} else {
return new ArrayList<FieldMeta>();
}
}
/**
* 处理desc table的结果
*/
public static List<FieldMeta> parseTableMetaByDesc(ResultSetPacket packet) {
Map<String, Integer> nameMaps = new HashMap<String, Integer>(6, 1f);
int index = 0;
for (FieldPacket fieldPacket : packet.getFieldDescriptors()) {
nameMaps.put(fieldPacket.getOriginalName(), index++);
@@ -103,7 +128,10 @@ public class TableMetaCache {
* size),
"YES"));
meta.setKey("PRI".equalsIgnoreCase(packet.getFieldValues().get(nameMaps.get(COLUMN_KEY) + i * size)));
meta.setDefaultValue(packet.getFieldValues().get(nameMaps.get(COLUMN_DEFAULT) + i * size));
meta.setUnique("UNI".equalsIgnoreCase(packet.getFieldValues().get(nameMaps.get(COLUMN_KEY) + i * size)));
// 特殊处理引号
meta.setDefaultValue(DruidDdlParser.unescapeQuotaName(packet.getFieldValues()
.get(nameMaps.get(COLUMN_DEFAULT) + i * size)));
meta.setExtra(packet.getFieldValues().get(nameMaps.get(EXTRA) + i * size));
result.add(meta);
@@ -4,35 +4,35 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import com.alibaba.druid.sql.SQLUtils;
import com.alibaba.druid.sql.ast.SQLExpr;
import com.alibaba.druid.sql.ast.SQLStatement;
import com.alibaba.druid.sql.ast.expr.SQLIdentifierExpr;
import com.alibaba.druid.sql.ast.expr.SQLPropertyExpr;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableAddConstraint;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableAddIndex;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableDropConstraint;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableDropIndex;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableDropKey;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableItem;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableRename;
import com.alibaba.druid.sql.ast.statement.SQLAlterTableStatement;
import com.alibaba.druid.sql.ast.statement.SQLConstraint;
import com.alibaba.druid.sql.ast.statement.SQLCreateIndexStatement;
import com.alibaba.druid.sql.ast.statement.SQLCreateTableStatement;
import com.alibaba.druid.sql.ast.statement.SQLDeleteStatement;
import com.alibaba.druid.sql.ast.statement.SQLDropIndexStatement;
import com.alibaba.druid.sql.ast.statement.SQLDropTableStatement;
import com.alibaba.druid.sql.ast.statement.SQLExprTableSource;
import com.alibaba.druid.sql.ast.statement.SQLInsertStatement;
import com.alibaba.druid.sql.ast.statement.SQLTableSource;
import com.alibaba.druid.sql.ast.statement.SQLTruncateStatement;
import com.alibaba.druid.sql.ast.statement.SQLUnique;
import com.alibaba.druid.sql.ast.statement.SQLUpdateStatement;
import com.alibaba.druid.sql.dialect.mysql.ast.statement.MySqlRenameTableStatement;
import com.alibaba.druid.sql.dialect.mysql.ast.statement.MySqlRenameTableStatement.Item;
import com.alibaba.druid.sql.parser.ParserException;
import com.alibaba.druid.util.JdbcConstants;
import com.alibaba.fastsql.sql.SQLUtils;
import com.alibaba.fastsql.sql.ast.SQLExpr;
import com.alibaba.fastsql.sql.ast.SQLStatement;
import com.alibaba.fastsql.sql.ast.expr.SQLIdentifierExpr;
import com.alibaba.fastsql.sql.ast.expr.SQLPropertyExpr;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableAddConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableAddIndex;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableDropConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableDropIndex;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableDropKey;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableItem;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableRename;
import com.alibaba.fastsql.sql.ast.statement.SQLAlterTableStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLCreateIndexStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLCreateTableStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLDeleteStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLDropIndexStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLDropTableStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLExprTableSource;
import com.alibaba.fastsql.sql.ast.statement.SQLInsertStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLTableSource;
import com.alibaba.fastsql.sql.ast.statement.SQLTruncateStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLUnique;
import com.alibaba.fastsql.sql.ast.statement.SQLUpdateStatement;
import com.alibaba.fastsql.sql.dialect.mysql.ast.statement.MySqlRenameTableStatement;
import com.alibaba.fastsql.sql.dialect.mysql.ast.statement.MySqlRenameTableStatement.Item;
import com.alibaba.fastsql.sql.parser.ParserException;
import com.alibaba.fastsql.util.JdbcConstants;
import com.alibaba.otter.canal.protocol.CanalEntry.EventType;
/**
@@ -187,7 +187,7 @@ public class DruidDdlParser {
}
public static String unescapeName(String name) {
if (name.length() > 2) {
if (name != null && name.length() > 2) {
char c0 = name.charAt(0);
char x0 = name.charAt(name.length() - 1);
if ((c0 == '"' && x0 == '"') || (c0 == '`' && x0 == '`')) {
@@ -198,4 +198,16 @@ public class DruidDdlParser {
return name;
}
public static String unescapeQuotaName(String name) {
if (name != null && name.length() > 2) {
char c0 = name.charAt(0);
char x0 = name.charAt(name.length() - 1);
if (c0 == '\'' && x0 == '\'') {
return name.substring(1, name.length() - 1);
}
}
return name;
}
}
@@ -18,9 +18,9 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import com.alibaba.druid.sql.repository.Schema;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.alibaba.fastsql.sql.repository.Schema;
import com.alibaba.otter.canal.filter.CanalEventFilter;
import com.alibaba.otter.canal.parse.driver.mysql.packets.server.ResultSetPacket;
import com.alibaba.otter.canal.parse.exception.CanalParseException;
@@ -44,19 +44,19 @@ import com.alibaba.otter.canal.protocol.position.EntryPosition;
*/
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 final EntryPosition INIT_POSITION = new EntryPosition("0", 0L, -2L, -1L);
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;
public static final EntryPosition INIT_POSITION = new EntryPosition("0", 0L, -2L, -1L);
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 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;
public DatabaseTableMeta(){
@@ -248,7 +248,7 @@ public class DatabaseTableMeta implements TableMetaTSDB {
boolean compareAll = true;
for (Schema schema : tmpMemoryTableMeta.getRepository().getSchemas()) {
for (String table : schema.showTables()) {
if (!compareTableMetaDbAndMemory(connection, schema.getName(), table)) {
if (!compareTableMetaDbAndMemory(connection, tmpMemoryTableMeta, schema.getName(), table)) {
compareAll = false;
}
}
@@ -284,15 +284,20 @@ public class DatabaseTableMeta implements TableMetaTSDB {
return false;
}
private boolean compareTableMetaDbAndMemory(MysqlConnection connection, final String schema, final String table) {
private boolean compareTableMetaDbAndMemory(MysqlConnection connection, MemoryTableMeta memoryTableMeta,
final String schema, final String table) {
TableMeta tableMetaFromMem = memoryTableMeta.find(schema, table);
TableMeta tableMetaFromDB = new TableMeta();
tableMetaFromDB.setSchema(schema);
tableMetaFromDB.setTable(table);
String createDDL = null;
try {
ResultSetPacket packet = connection.query("desc " + getFullName(schema, table));
tableMetaFromDB.setFields(TableMetaCache.parserTableMeta(packet));
ResultSetPacket packet = connection.query("show create table " + getFullName(schema, table));
if (packet.getFieldValues().size() > 1) {
createDDL = packet.getFieldValues().get(1);
tableMetaFromDB.setFields(TableMetaCache.parseTableMeta(schema, table, packet));
}
} catch (IOException e) {
if (e.getMessage().contains("errorNumber=1146")) {
logger.error("table not exist in db , pls check :" + getFullName(schema, table) + " , mem : "
@@ -304,16 +309,6 @@ public class DatabaseTableMeta implements TableMetaTSDB {
boolean result = compareTableMeta(tableMetaFromMem, tableMetaFromDB);
if (!result) {
String createDDL = null;
try {
ResultSetPacket packet = connection.query("show create table " + getFullName(schema, table));
if (packet.getFieldValues().size() > 1) {
createDDL = packet.getFieldValues().get(1);
}
} catch (IOException e) {
// ignore
}
logger.error("pls submit github issue, show create table ddl:" + createDDL + " , compare failed . \n db : "
+ tableMetaFromDB + " \n mem : " + tableMetaFromMem);
}
@@ -442,7 +437,10 @@ public class DatabaseTableMeta implements TableMetaTSDB {
return false;
}
if (sourceField.isKey() != targetField.isKey()) {
// mysql会有一种处理,针对show create只有uk没有pk时,会在desc默认将uk当做pk
boolean isSourcePkOrUk = sourceField.isKey() || sourceField.isUnique();
boolean isTargetPkOrUk = targetField.isKey() || targetField.isUnique();
if (isSourcePkOrUk != isTargetPkOrUk) {
return false;
}
}
@@ -10,27 +10,29 @@ import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.alibaba.druid.sql.ast.SQLDataType;
import com.alibaba.druid.sql.ast.SQLDataTypeImpl;
import com.alibaba.druid.sql.ast.SQLExpr;
import com.alibaba.druid.sql.ast.SQLStatement;
import com.alibaba.druid.sql.ast.expr.SQLCharExpr;
import com.alibaba.druid.sql.ast.expr.SQLIdentifierExpr;
import com.alibaba.druid.sql.ast.expr.SQLNullExpr;
import com.alibaba.druid.sql.ast.expr.SQLPropertyExpr;
import com.alibaba.druid.sql.ast.statement.SQLColumnConstraint;
import com.alibaba.druid.sql.ast.statement.SQLColumnDefinition;
import com.alibaba.druid.sql.ast.statement.SQLColumnPrimaryKey;
import com.alibaba.druid.sql.ast.statement.SQLCreateTableStatement;
import com.alibaba.druid.sql.ast.statement.SQLNotNullConstraint;
import com.alibaba.druid.sql.ast.statement.SQLNullConstraint;
import com.alibaba.druid.sql.ast.statement.SQLSelectOrderByItem;
import com.alibaba.druid.sql.ast.statement.SQLTableElement;
import com.alibaba.druid.sql.dialect.mysql.ast.MySqlPrimaryKey;
import com.alibaba.druid.sql.repository.Schema;
import com.alibaba.druid.sql.repository.SchemaObject;
import com.alibaba.druid.sql.repository.SchemaRepository;
import com.alibaba.druid.util.JdbcConstants;
import com.alibaba.fastsql.sql.ast.SQLDataType;
import com.alibaba.fastsql.sql.ast.SQLDataTypeImpl;
import com.alibaba.fastsql.sql.ast.SQLExpr;
import com.alibaba.fastsql.sql.ast.SQLStatement;
import com.alibaba.fastsql.sql.ast.expr.SQLCharExpr;
import com.alibaba.fastsql.sql.ast.expr.SQLIdentifierExpr;
import com.alibaba.fastsql.sql.ast.expr.SQLNullExpr;
import com.alibaba.fastsql.sql.ast.expr.SQLPropertyExpr;
import com.alibaba.fastsql.sql.ast.statement.SQLColumnConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLColumnDefinition;
import com.alibaba.fastsql.sql.ast.statement.SQLColumnPrimaryKey;
import com.alibaba.fastsql.sql.ast.statement.SQLColumnUniqueKey;
import com.alibaba.fastsql.sql.ast.statement.SQLCreateTableStatement;
import com.alibaba.fastsql.sql.ast.statement.SQLNotNullConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLNullConstraint;
import com.alibaba.fastsql.sql.ast.statement.SQLSelectOrderByItem;
import com.alibaba.fastsql.sql.ast.statement.SQLTableElement;
import com.alibaba.fastsql.sql.dialect.mysql.ast.MySqlPrimaryKey;
import com.alibaba.fastsql.sql.dialect.mysql.ast.MySqlUnique;
import com.alibaba.fastsql.sql.repository.Schema;
import com.alibaba.fastsql.sql.repository.SchemaObject;
import com.alibaba.fastsql.sql.repository.SchemaRepository;
import com.alibaba.fastsql.util.JdbcConstants;
import com.alibaba.otter.canal.parse.inbound.TableMeta;
import com.alibaba.otter.canal.parse.inbound.TableMeta.FieldMeta;
import com.alibaba.otter.canal.parse.inbound.mysql.ddl.DruidDdlParser;
@@ -64,7 +66,10 @@ public class MemoryTableMeta implements TableMetaTSDB {
}
try {
repository.console(ddl);
// druid暂时flush privileges语法解析有问题
if (!StringUtils.startsWithIgnoreCase(StringUtils.trim(ddl), "flush")) {
repository.console(ddl);
}
} catch (Throwable e) {
logger.warn("parse faield : " + ddl, e);
}
@@ -187,7 +192,7 @@ public class MemoryTableMeta implements TableMetaTSDB {
if (column.getDefaultExpr() == null || column.getDefaultExpr() instanceof SQLNullExpr) {
fieldMeta.setDefaultValue(null);
} else {
fieldMeta.setDefaultValue(getSqlName(column.getDefaultExpr()));
fieldMeta.setDefaultValue(DruidDdlParser.unescapeQuotaName(getSqlName(column.getDefaultExpr())));
}
fieldMeta.setColumnName(name);
@@ -201,6 +206,8 @@ public class MemoryTableMeta implements TableMetaTSDB {
fieldMeta.setNullable(true);
} else if (constraint instanceof SQLColumnPrimaryKey) {
fieldMeta.setKey(true);
} else if (constraint instanceof SQLColumnUniqueKey) {
fieldMeta.setUnique(true);
}
}
tableMeta.addFieldMeta(fieldMeta);
@@ -212,6 +219,14 @@ public class MemoryTableMeta implements TableMetaTSDB {
FieldMeta field = tableMeta.getFieldMetaByName(name);
field.setKey(true);
}
} else if (element instanceof MySqlUnique) {
MySqlUnique column = (MySqlUnique) element;
List<SQLSelectOrderByItem> uks = column.getColumns();
for (SQLSelectOrderByItem uk : uks) {
String name = getSqlName(uk.getExpr());
FieldMeta field = tableMeta.getFieldMetaByName(name);
field.setUnique(true);
}
}
}
@@ -3,27 +3,38 @@ package com.alibaba.otter.canal.parse;
import java.io.IOException;
import java.net.InetSocketAddress;
import org.apache.commons.lang.StringUtils;
import org.junit.Assert;
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.alibaba.otter.canal.parse.driver.mysql.MysqlConnector;
import com.alibaba.otter.canal.parse.driver.mysql.MysqlUpdateExecutor;
import com.alibaba.otter.canal.parse.driver.mysql.packets.HeaderPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.client.BinlogDumpCommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.client.RegisterSlaveCommandPacket;
import com.alibaba.otter.canal.parse.driver.mysql.packets.server.ErrorPacket;
import com.alibaba.otter.canal.parse.driver.mysql.utils.PacketManager;
import com.alibaba.otter.canal.parse.inbound.mysql.dbsync.DirectLogFetcher;
import com.taobao.tddl.dbsync.binlog.LogContext;
import com.taobao.tddl.dbsync.binlog.LogDecoder;
import com.taobao.tddl.dbsync.binlog.LogEvent;
import com.taobao.tddl.dbsync.binlog.event.RotateLogEvent;
public class DirectLogFetcherTest {
protected final Logger logger = LoggerFactory.getLogger(this.getClass());
@Test
public void testSimple() {
DirectLogFetcher fetcher = new DirectLogFetcher();
try {
MysqlConnector connector = new MysqlConnector(new InetSocketAddress("127.0.0.1", 3306), "xxxxx", "xxxxx");
MysqlConnector connector = new MysqlConnector(new InetSocketAddress("127.0.0.1", 3306), "xxxx", "xxxx");
connector.connect();
sendBinlogDump(connector, "mysql-bin.001016", 4L, 3);
updateSettings(connector);
sendRegisterSlave(connector, 3);
sendBinlogDump(connector, "mysql-bin.000001", 4L, 3);
fetcher.start(connector.getChannel());
@@ -42,6 +53,7 @@ public class DirectLogFetcherTest {
case LogEvent.ROTATE_EVENT:
// binlogFileName = ((RotateLogEvent)
// event).getFilename();
System.out.println(((RotateLogEvent) event).getFilename());
break;
case LogEvent.WRITE_ROWS_EVENT_V1:
case LogEvent.WRITE_ROWS_EVENT:
@@ -82,6 +94,33 @@ public class DirectLogFetcherTest {
}
private void sendRegisterSlave(MysqlConnector connector, int slaveId) throws IOException {
RegisterSlaveCommandPacket cmd = new RegisterSlaveCommandPacket();
cmd.reportHost = connector.getAddress().getAddress().getHostAddress();
cmd.reportPasswd = connector.getPassword();
cmd.reportUser = connector.getUsername();
cmd.serverId = slaveId;
byte[] cmdBody = cmd.toBytes();
HeaderPacket header = new HeaderPacket();
header.setPacketBodyLength(cmdBody.length);
header.setPacketSequenceNumber((byte) 0x00);
PacketManager.writePkg(connector.getChannel(), header.toBytes(), cmdBody);
header = PacketManager.readHeader(connector.getChannel(), 4);
byte[] body = PacketManager.readBytes(connector.getChannel(), header.getPacketBodyLength());
assert body != null;
if (body[0] < 0) {
if (body[0] == -1) {
ErrorPacket err = new ErrorPacket();
err.fromBytes(body);
throw new IOException("Error When doing Register slave:" + err.toString());
} else {
throw new IOException("unpexpected packet with field_count=" + body[0]);
}
}
}
private void sendBinlogDump(MysqlConnector connector, String binlogfilename, Long binlogPosition, int slaveId)
throws IOException {
BinlogDumpCommandPacket binlogDumpCmd = new BinlogDumpCommandPacket();
@@ -95,4 +134,63 @@ public class DirectLogFetcherTest {
binlogDumpHeader.setPacketSequenceNumber((byte) 0x00);
PacketManager.writePkg(connector.getChannel(), binlogDumpHeader.toBytes(), cmdBody);
}
private void updateSettings(MysqlConnector connector) throws IOException {
try {
update("set wait_timeout=9999999", connector);
} catch (Exception e) {
logger.warn("update wait_timeout failed", e);
}
try {
update("set net_write_timeout=1800", connector);
} catch (Exception e) {
logger.warn("update net_write_timeout failed", e);
}
try {
update("set net_read_timeout=1800", connector);
} catch (Exception e) {
logger.warn("update net_read_timeout failed", e);
}
try {
// 设置服务端返回结果时不做编码转化,直接按照数据库的二进制编码进行发送,由客户端自己根据需求进行编码转化
update("set names 'binary'", connector);
} catch (Exception e) {
logger.warn("update names failed", e);
}
try {
// mysql5.6针对checksum支持需要设置session变量
// 如果不设置会出现错误: Slave can not handle replication events with the
// checksum that master is configured to log
// 但也不能乱设置,需要和mysql server的checksum配置一致,不然RotateLogEvent会出现乱码
update("set @master_binlog_checksum= '@@global.binlog_checksum'", connector);
} catch (Exception e) {
logger.warn("update master_binlog_checksum failed", e);
}
try {
// 参考:https://github.com/alibaba/canal/issues/284
// mysql5.6需要设置slave_uuid避免被server kill链接
update("set @slave_uuid=uuid()", connector);
} catch (Exception e) {
if (!StringUtils.contains(e.getMessage(), "Unknown system variable")) {
logger.warn("update slave_uuid failed", e);
}
}
try {
// mariadb针对特殊的类型,需要设置session变量
update("SET @mariadb_slave_capability='" + LogEvent.MARIA_SLAVE_CAPABILITY_MINE + "'", connector);
} catch (Exception e) {
logger.warn("update mariadb_slave_capability failed", e);
}
}
public void update(String cmd, MysqlConnector connector) throws IOException {
MysqlUpdateExecutor exector = new MysqlUpdateExecutor(connector);
exector.update(cmd);
}
}
+64 -23
View File
@@ -1,23 +1,64 @@
CREATE TABLE `test` (
`id` bigint(20) zerofill unsigNed NOT NULL AUTO_INCREMENT COMMENT 'id',
`c_tinyint` tinyint(4) DEFAULT '1' COMMENT 'tinyint',
`c_smallint` smallint(6) DEFAULT 0 COMMENT 'smallint',
`c_mediumint` mediumint(9) DEFAULT NULL COMMENT 'mediumint',
`c_int` int(11) DEFAULT NULL COMMENT 'int',
`c_bigint` bigint(20) DEFAULT NULL COMMENT 'bigint',
`c_decimal` decimal(10,3) DEFAULT NULL COMMENT 'decimal',
`c_date` date DEFAULT '0000-00-00' COMMENT 'date',
`c_datetime` datetime DEFAULT '0000-00-00 00:00:00' COMMENT 'datetime',
`c_timestamp` timestamp NULL DEFAULT NULL COMMENT 'timestamp',
`c_time` time DEFAULT NULL COMMENT 'time',
`c_char` char(10) DEFAULT NULL COMMENT 'char',
`c_varchar` varchar(10) DEFAULT 'hello' COMMENT 'varchar',
`c_blob` blob COMMENT 'blob',
`c_text` text COMMENT 'text',
`c_mediumtext` mediumtext COMMENT 'mediumtext',
`c_longblob` longblob COMMENT 'longblob',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_a` (`c_tinyint`),
KEY `k_b` (`c_smallint`),
KEY `k_c` (`c_mediumint`,`c_int`)
) ENGINE=InnoDB AUTO_INCREMENT=1769503 DEFAULT CHARSET=utf8mb4 COMMENT='10000000';
CREATE TABLE `test_all` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`c_bit_1` bit(1) DEFAULT NULL,
`c_bit_8` bit(8) DEFAULT NULL,
`c_bit_16` bit(16) DEFAULT NULL,
`c_bit_32` bit(32) DEFAULT NULL,
`c_bit_64` bit(64) DEFAULT NULL,
`c_bool` boolean DEFAULT NULL,
`c_tinyint_1` tinyint(1) DEFAULT NULL,
`c_tinyint_4` tinyint(4) DEFAULT NULL,
`c_tinyint_8` tinyint(8) DEFAULT NULL,
`c_tinyint_8_un` tinyint(8) unsigned DEFAULT NULL,
`c_smallint_1` smallint(1) DEFAULT NULL,
`c_smallint_16` smallint(16) DEFAULT NULL,
`c_smallint_16_un` smallint(16) unsigned DEFAULT NULL,
`c_mediumint_1` mediumint(1) DEFAULT NULL,
`c_mediumint_24` mediumint(24) DEFAULT NULL,
`c_mediumint_24_un` mediumint(24) unsigned DEFAULT NULL,
`c_int_1` int(1) DEFAULT NULL,
`c_int_32` int(32) DEFAULT NULL,
`c_int_32_un` int(32) unsigned DEFAULT NULL,
`c_bigint_1` bigint(1) DEFAULT NULL,
`c_bigint_64` bigint(64) DEFAULT NULL,
`c_bigint_64_un` bigint(64) unsigned DEFAULT NULL,
`c_decimal` decimal(10,3) DEFAULT NULL,
`c_decimal_pr` decimal(10,3) DEFAULT NULL,
`c_float` float DEFAULT NULL,
`c_float_pr` float(10,3) DEFAULT NULL,
`c_float_un` float(10,3) unsigned DEFAULT NULL,
`c_double` double DEFAULT NULL,
`c_double_pr` double(10,3) DEFAULT NULL,
`c_double_un` double(10,3) unsigned DEFAULT NULL,
`c_date` date DEFAULT NULL COMMENT 'date',
`c_datetime` datetime DEFAULT NULL,
`c_datetime_1` datetime(1) DEFAULT NULL,
`c_datetime_3` datetime(3) DEFAULT NULL,
`c_datetime_6` datetime(6) DEFAULT NULL,
`c_timestamp` timestamp DEFAULT CURRENT_TIMESTAMP,
`c_timestamp_1` timestamp(1) DEFAULT 0,
`c_timestamp_3` timestamp(3) DEFAULT 0,
`c_timestamp_6` timestamp(6) DEFAULT 0,
`c_time` time DEFAULT NULL,
`c_time_1` time(1) DEFAULT NULL,
`c_time_3` time(3) DEFAULT NULL,
`c_time_6` time(6) DEFAULT NULL,
`c_year` year DEFAULT NULL,
`c_year_4` year(4) DEFAULT NULL,
`c_char` char(10) DEFAULT NULL,
`c_varchar` varchar(10) DEFAULT NULL,
`c_binary` binary(10) DEFAULT NULL,
`c_varbinary` varbinary(10) DEFAULT NULL,
`c_blob_tiny` tinyblob DEFAULT NULL,
`c_blob` blob DEFAULT NULL,
`c_blob_medium` mediumblob DEFAULT NULL,
`c_blob_long` longblob DEFAULT NULL,
`c_text_tiny` tinytext DEFAULT NULL,
`c_text` text DEFAULT NULL,
`c_text_medium` mediumtext DEFAULT NULL,
`c_text_long` longtext DEFAULT NULL,
`c_enum` enum('a','b','c') DEFAULT NULL,
`c_set` set('a','b','c') DEFAULT NULL,
`c_json` json DEFAULT NULL,
PRIMARY KEY (`id`)
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8mb4 COMMENT='10000000'
+15 -2
View File
@@ -4,7 +4,7 @@
<artifactId>canal</artifactId>
<packaging>pom</packaging>
<name>canal module for otter ${project.version}</name>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<url>https://github.com/alibaba/canal</url>
<parent>
<groupId>org.sonatype.oss</groupId>
@@ -251,10 +251,15 @@
<artifactId>ibatis-sqlmap</artifactId>
<version>2.3.4.726</version>
</dependency>
<dependency>
<groupId>com.alibaba.fastsql</groupId>
<artifactId>fastsql</artifactId>
<version>2.0.0_preview_159</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>druid</artifactId>
<version>1.1.5-preview_14</version>
<version>1.1.9</version>
</dependency>
<!-- log -->
<dependency>
@@ -445,6 +450,14 @@
</excludes>
</testResource>
</testResources>
<pluginManagement>
<plugins>
<plugin>
<artifactId>maven-jar-plugin</artifactId>
<version>3.0.2</version>
</plugin>
</plugins>
</pluginManagement>
</build>
<distributionManagement>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<artifactId>canal.server</artifactId>
@@ -220,8 +220,8 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C
events = getEvents(canalInstance.getEventStore(), start, batchSize, timeout, unit);
if (CollectionUtils.isEmpty(events.getEvents())) {
logger.debug("get successfully, clientId:{} batchSize:{} but result is null", new Object[] {
clientIdentity.getClientId(), batchSize });
logger.debug("get successfully, clientId:{} batchSize:{} but result is null",
clientIdentity.getClientId(), batchSize);
return new Message(-1, new ArrayList<Entry>()); // 返回空包,避免生成batchId,浪费性能
} else {
// 记录到流式信息
@@ -232,13 +232,14 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C
return input.getEntry();
}
});
logger.info("get successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]",
clientIdentity.getClientId(),
batchSize,
entrys.size(),
batchId,
events.getPositionRange());
if (logger.isInfoEnabled()) {
logger.info("get successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]",
clientIdentity.getClientId(),
batchSize,
entrys.size(),
batchId,
events.getPositionRange());
}
// 直接提交ack
ack(clientIdentity, batchId);
return new Message(batchId, entrys);
@@ -297,8 +298,8 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C
}
if (CollectionUtils.isEmpty(events.getEvents())) {
logger.debug("getWithoutAck successfully, clientId:{} batchSize:{} but result is null", new Object[] {
clientIdentity.getClientId(), batchSize });
logger.debug("getWithoutAck successfully, clientId:{} batchSize:{} but result is null",
clientIdentity.getClientId(), batchSize);
return new Message(-1, new ArrayList<Entry>()); // 返回空包,避免生成batchId,浪费性能
} else {
// 记录到流式信息
@@ -309,13 +310,14 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C
return input.getEntry();
}
});
logger.info("getWithoutAck successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]",
clientIdentity.getClientId(),
batchSize,
entrys.size(),
batchId,
events.getPositionRange());
if (logger.isInfoEnabled()) {
logger.info("getWithoutAck successfully, clientId:{} batchSize:{} real size is {} and result is [batchId:{} , position:{}]",
clientIdentity.getClientId(),
batchSize,
entrys.size(),
batchId,
events.getPositionRange());
}
return new Message(batchId, entrys);
}
@@ -377,10 +379,12 @@ public class CanalServerWithEmbedded extends AbstractCanalLifeCycle implements C
// 更新cursor
if (positionRanges.getAck() != null) {
canalInstance.getMetaManager().updateCursor(clientIdentity, positionRanges.getAck());
logger.info("ack successfully, clientId:{} batchId:{} position:{}",
clientIdentity.getClientId(),
batchId,
positionRanges);
if (logger.isInfoEnabled()) {
logger.info("ack successfully, clientId:{} batchId:{} position:{}",
clientIdentity.getClientId(),
batchId,
positionRanges);
}
}
// 可定时清理数据
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>
@@ -5,6 +5,7 @@ import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import com.alibaba.otter.canal.protocol.CanalEntry.EntryType;
import com.alibaba.otter.canal.sink.exception.CanalSinkException;
import com.alibaba.otter.canal.store.model.Event;
/**
@@ -59,21 +60,22 @@ public class TimelineTransactionBarrier extends TimelineBarrier {
}
public void clear(Event event) {
super.clear(event);
super.clear(event);
if (isTransactionEnd(event)) {
//应该先判断2,再判断是否是事务尾,因为事务尾也可以导致txState的状态为2
//如果先判断事务尾,那么2的状态可能永远没机会被修改了,系统出现死锁
//CanalSinkException被注释的代码是不是可以放开??我们内部使用的时候已经放开了,从代码逻辑的分析上以及实践效果来看,应该抛异常
if (txState.intValue() == 2) {// 非事务中
boolean result = txState.compareAndSet(2, 0);
if (result == false) {
throw new CanalSinkException("state is not correct in non-transaction");
}
} else if (isTransactionEnd(event)) {
inTransaction.set(false); // 事务结束并且已经成功写入store,清理标记,进入重新排队判断,允许新的事务进入
txState.compareAndSet(1, 0);
// if (txState.compareAndSet(1, 0) == false) {
// throw new
// CanalSinkException("state is not correct in transaction");
// }
} else if (txState.intValue() == 2) {// 非事务中
txState.compareAndSet(2, 0);
// if (txState.compareAndSet(2, 0) == false) {
// throw new
// CanalSinkException("state is not correct in non-transaction");
// }
boolean result = txState.compareAndSet(1, 0);
if (result == false) {
throw new CanalSinkException("state is not correct in transaction");
}
}
}
@@ -90,7 +92,8 @@ public class TimelineTransactionBarrier extends TimelineBarrier {
return true; // 事务允许通过
}
} else if (txState.compareAndSet(0, 2)) { // 非事务保护中
return true; // DDL/DCL允许通过
//当基于zk-cursor启动的时候,拿到的第一个Event是TransactionEnd
return true; // DDL/DCL/TransactionEnd允许通过
}
}
}
+1 -1
View File
@@ -3,7 +3,7 @@
<parent>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal</artifactId>
<version>1.0.25-SNAPSHOT</version>
<version>1.0.26-SNAPSHOT</version>
<relativePath>../pom.xml</relativePath>
</parent>
<groupId>com.alibaba.otter</groupId>