>();
+ private static final ConvertUtilsBean convertUtilsBean = new ConvertUtilsBean();
+
+ static {
+ // regist Converter
+ convertUtilsBean.register(SqlTimestampConverter.SQL_TIMESTAMP, Date.class);
+ convertUtilsBean.register(SqlTimestampConverter.SQL_TIMESTAMP, Time.class);
+ convertUtilsBean.register(SqlTimestampConverter.SQL_TIMESTAMP, Timestamp.class);
+ convertUtilsBean.register(ByteArrayConverter.SQL_BYTES, byte[].class);
+
+ // bool
+ sqlTypeToJavaTypeMap.put(Types.BOOLEAN, Boolean.class);
+
+ // int
+ sqlTypeToJavaTypeMap.put(Types.TINYINT, Integer.class);
+ sqlTypeToJavaTypeMap.put(Types.SMALLINT, Integer.class);
+ sqlTypeToJavaTypeMap.put(Types.INTEGER, Integer.class);
+
+ // long
+ sqlTypeToJavaTypeMap.put(Types.BIGINT, Long.class);
+ // mysql bit最多64位,无符号
+ sqlTypeToJavaTypeMap.put(Types.BIT, BigInteger.class);
+
+ // decimal
+ sqlTypeToJavaTypeMap.put(Types.REAL, Float.class);
+ sqlTypeToJavaTypeMap.put(Types.FLOAT, Float.class);
+ sqlTypeToJavaTypeMap.put(Types.DOUBLE, Double.class);
+ sqlTypeToJavaTypeMap.put(Types.NUMERIC, BigDecimal.class);
+ sqlTypeToJavaTypeMap.put(Types.DECIMAL, BigDecimal.class);
+
+ // date
+ sqlTypeToJavaTypeMap.put(Types.DATE, Date.class);
+ sqlTypeToJavaTypeMap.put(Types.TIME, Time.class);
+ sqlTypeToJavaTypeMap.put(Types.TIMESTAMP, Timestamp.class);
+
+ // blob
+ sqlTypeToJavaTypeMap.put(Types.BLOB, byte[].class);
+
+ // byte[]
+ sqlTypeToJavaTypeMap.put(Types.REF, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.OTHER, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.ARRAY, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.STRUCT, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.SQLXML, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.BINARY, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.DATALINK, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.DISTINCT, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.VARBINARY, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.JAVA_OBJECT, byte[].class);
+ sqlTypeToJavaTypeMap.put(Types.LONGVARBINARY, byte[].class);
+
+ // String
+ sqlTypeToJavaTypeMap.put(Types.CHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.VARCHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.LONGVARCHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.LONGNVARCHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.NCHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.NVARCHAR, String.class);
+ sqlTypeToJavaTypeMap.put(Types.NCLOB, String.class);
+ sqlTypeToJavaTypeMap.put(Types.CLOB, String.class);
+ }
+
+ /**
+ * 将指定java.sql.Types的ResultSet value转换成相应的String
+ *
+ * @param rs
+ * @param index
+ * @param sqlType
+ * @return
+ * @throws SQLException
+ */
+ public static String sqlValueToString(ResultSet rs, int index, int sqlType) throws SQLException {
+ Class> requiredType = sqlTypeToJavaTypeMap.get(sqlType);
+ if (requiredType == null) {
+ throw new IllegalArgumentException("unknow java.sql.Types - " + sqlType);
+ }
+
+ return getResultSetValue(rs, index, requiredType);
+ }
+
+ /**
+ * sqlValueToString方法的逆向过程
+ *
+ * @param value
+ * @param sqlType
+ * @param isTextRequired
+ * @param isEmptyStringNulled
+ * @return
+ */
+ public static Object stringToSqlValue(String value, int sqlType, boolean isRequired, boolean isEmptyStringNulled) {
+ // 设置变量
+ String sourceValue = value;
+ if (SqlUtils.isTextType(sqlType)) {
+ if ((sourceValue == null) || (StringUtils.isEmpty(sourceValue) && isEmptyStringNulled)) {
+ return isRequired ? REQUIRED_FIELD_NULL_SUBSTITUTE : null;
+ } else {
+ return sourceValue;
+ }
+ } else {
+ if (StringUtils.isEmpty(sourceValue)) {
+ return null;
+ } else {
+ Class> requiredType = sqlTypeToJavaTypeMap.get(sqlType);
+ if (requiredType == null) {
+ throw new IllegalArgumentException("unknow java.sql.Types - " + sqlType);
+ } else if (requiredType.equals(String.class)) {
+ return sourceValue;
+ } else if (isNumeric(sqlType)) {
+ return convertUtilsBean.convert(sourceValue.trim(), requiredType);
+ } else {
+ return convertUtilsBean.convert(sourceValue, requiredType);
+ }
+ }
+ }
+ }
+
+ public static String encoding(String source, int sqlType, String sourceEncoding, String targetEncoding) {
+ switch (sqlType) {
+ case Types.CHAR:
+ case Types.VARCHAR:
+ case Types.LONGVARCHAR:
+ case Types.NCHAR:
+ case Types.NVARCHAR:
+ case Types.LONGNVARCHAR:
+ case Types.CLOB:
+ case Types.NCLOB:
+ if (false == StringUtils.isEmpty(source)) {
+ String fromEncoding = StringUtils.isBlank(sourceEncoding) ? "UTF-8" : sourceEncoding;
+ String toEncoding = StringUtils.isBlank(targetEncoding) ? "UTF-8" : targetEncoding;
+
+ // if (false == StringUtils.equalsIgnoreCase(fromEncoding,
+ // toEncoding)) {
+ try {
+ return new String(source.getBytes(fromEncoding), toEncoding);
+ } catch (UnsupportedEncodingException e) {
+ throw new IllegalArgumentException(e.getMessage(), e);
+ }
+ // }
+ }
+ }
+
+ return source;
+ }
+
+ /**
+ * Retrieve a JDBC column value from a ResultSet, using the specified value
+ * type.
+ *
+ * Uses the specifically typed ResultSet accessor methods, falling back to
+ * {@link #getResultSetValue(ResultSet, int)} for unknown types.
+ *
+ * Note that the returned value may not be assignable to the specified
+ * required type, in case of an unknown type. Calling code needs to deal
+ * with this case appropriately, e.g. throwing a corresponding exception.
+ *
+ * @param rs is the ResultSet holding the data
+ * @param index is the column index
+ * @param requiredType the required value type (may be null)
+ * @return the value object
+ * @throws SQLException if thrown by the JDBC API
+ */
+ private static String getResultSetValue(ResultSet rs, int index, Class> requiredType) throws SQLException {
+ if (requiredType == null) {
+ return getResultSetValue(rs, index);
+ }
+
+ Object value = null;
+ boolean wasNullCheck = false;
+
+ // Explicitly extract typed value, as far as possible.
+ if (String.class.equals(requiredType)) {
+ value = rs.getString(index);
+ } else if (boolean.class.equals(requiredType) || Boolean.class.equals(requiredType)) {
+ value = Boolean.valueOf(rs.getBoolean(index));
+ wasNullCheck = true;
+ } else if (byte.class.equals(requiredType) || Byte.class.equals(requiredType)) {
+ value = new Byte(rs.getByte(index));
+ wasNullCheck = true;
+ } else if (short.class.equals(requiredType) || Short.class.equals(requiredType)) {
+ value = new Short(rs.getShort(index));
+ wasNullCheck = true;
+ } else if (int.class.equals(requiredType) || Integer.class.equals(requiredType)) {
+ value = new Long(rs.getLong(index));
+ wasNullCheck = true;
+ } else if (long.class.equals(requiredType) || Long.class.equals(requiredType)) {
+ value = rs.getBigDecimal(index);
+ wasNullCheck = true;
+ } else if (float.class.equals(requiredType) || Float.class.equals(requiredType)) {
+ value = new Float(rs.getFloat(index));
+ wasNullCheck = true;
+ } else if (double.class.equals(requiredType) || Double.class.equals(requiredType)
+ || Number.class.equals(requiredType)) {
+ value = new Double(rs.getDouble(index));
+ wasNullCheck = true;
+ } else if (Time.class.equals(requiredType)) {
+ // try {
+ // value = rs.getTime(index);
+ // } catch (SQLException e) {
+ value = rs.getString(index);// 尝试拿为string对象,0000无法用Time表示
+ // if (value == null && !rs.wasNull()) {
+ // value = "00:00:00"; //
+ // mysql设置了zeroDateTimeBehavior=convertToNull,出现0值时返回为null
+ // }
+ // }
+ } else if (Timestamp.class.equals(requiredType) || Date.class.equals(requiredType)) {
+ // try {
+ // value = convertTimestamp(rs.getTimestamp(index));
+ // } catch (SQLException e) {
+ // 尝试拿为string对象,0000-00-00 00:00:00无法用Timestamp 表示
+ value = rs.getString(index);
+ // if (value == null && !rs.wasNull()) {
+ // value = "0000:00:00 00:00:00"; //
+ // mysql设置了zeroDateTimeBehavior=convertToNull,出现0值时返回为null
+ // }
+ // }
+ } else if (BigDecimal.class.equals(requiredType)) {
+ value = rs.getBigDecimal(index);
+ } else if (BigInteger.class.equals(requiredType)) {
+ value = rs.getBigDecimal(index);
+ } else if (Blob.class.equals(requiredType)) {
+ value = rs.getBlob(index);
+ } else if (Clob.class.equals(requiredType)) {
+ value = rs.getClob(index);
+ } else if (byte[].class.equals(requiredType)) {
+ try {
+ byte[] bytes = rs.getBytes(index);
+ if (bytes == null) {
+ value = null;
+ } else {
+ value = new String(bytes, "ISO-8859-1");// 将binary转化为iso-8859-1的字符串
+ }
+ } catch (UnsupportedEncodingException e) {
+ throw new SQLException(e);
+ }
+ } else {
+ // Some unknown type desired -> rely on getObject.
+ value = getResultSetValue(rs, index);
+ }
+
+ // Perform was-null check if demanded (for results that the
+ // JDBC driver returns as primitives).
+ if (wasNullCheck && (value != null) && rs.wasNull()) {
+ value = null;
+ }
+
+ return (value == null) ? null : convertUtilsBean.convert(value);
+ }
+
+ /**
+ * Retrieve a JDBC column value from a ResultSet, using the most appropriate
+ * value type. The returned value should be a detached value object, not
+ * having any ties to the active ResultSet: in particular, it should not be
+ * a Blob or Clob object but rather a byte array respectively String
+ * representation.
+ *
+ * Uses the getObject(index) method, but includes additional
+ * "hacks" to get around Oracle 10g returning a non-standard object for its
+ * TIMESTAMP datatype and a java.sql.Date for DATE columns
+ * leaving out the time portion: These columns will explicitly be extracted
+ * as standard java.sql.Timestamp object.
+ *
+ * @param rs is the ResultSet holding the data
+ * @param index is the column index
+ * @return the value object
+ * @throws SQLException if thrown by the JDBC API
+ * @see Blob
+ * @see Clob
+ * @see Timestamp
+ */
+ private static String getResultSetValue(ResultSet rs, int index) throws SQLException {
+ Object obj = rs.getObject(index);
+ return (obj == null) ? null : convertUtilsBean.convert(obj);
+ }
+
+ // private static Object convertTimestamp(Timestamp timestamp) {
+ // return (timestamp == null) ? null : timestamp.getTime();
+ // }
+
+ /**
+ * Check whether the given SQL type is numeric.
+ */
+ public static boolean isNumeric(int sqlType) {
+ return (Types.BIT == sqlType) || (Types.BIGINT == sqlType) || (Types.DECIMAL == sqlType)
+ || (Types.DOUBLE == sqlType) || (Types.FLOAT == sqlType) || (Types.INTEGER == sqlType)
+ || (Types.NUMERIC == sqlType) || (Types.REAL == sqlType) || (Types.SMALLINT == sqlType)
+ || (Types.TINYINT == sqlType);
+ }
+
+ public static boolean isTextType(int sqlType) {
+ if (sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.CLOB || sqlType == Types.LONGVARCHAR
+ || sqlType == Types.NCHAR || sqlType == Types.NVARCHAR || sqlType == Types.NCLOB
+ || sqlType == Types.LONGNVARCHAR) {
+ return true;
+ } else {
+ return false;
+ }
+ }
+}
diff --git a/example/src/main/resources/client-spring.xml b/example/src/main/resources/client-spring.xml
new file mode 100644
index 00000000..29eba1af
--- /dev/null
+++ b/example/src/main/resources/client-spring.xml
@@ -0,0 +1,53 @@
+
+
+
+
+
+
+
+
+ classpath:client.properties
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/example/src/main/resources/client.properties b/example/src/main/resources/client.properties
new file mode 100644
index 00000000..c5142912
--- /dev/null
+++ b/example/src/main/resources/client.properties
@@ -0,0 +1,16 @@
+# client 配置
+zk.servers=127.0.0.1:2181
+# 5 * 1024
+client.batch.size=5120
+client.debug=false
+client.destination=example
+client.username=canal
+client.password=canal
+client.exceptionstrategy=1
+client.retrytimes=3
+client.filter=.*\\..*
+
+# 同步目标: mysql 配置
+target.mysql.url=jdbc:mysql://127.0.0.1:4306
+target.mysql.username=root
+target.mysql.password=123456
diff --git a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java
index 35e3235a..63ed2d02 100644
--- a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java
+++ b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/CanalInstanceWithManager.java
@@ -6,7 +6,6 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
-import com.alibaba.otter.canal.meta.FileMixedMetaManager;
import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -21,7 +20,14 @@ import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter;
import com.alibaba.otter.canal.instance.core.AbstractCanalInstance;
import com.alibaba.otter.canal.instance.manager.model.Canal;
import com.alibaba.otter.canal.instance.manager.model.CanalParameter;
-import com.alibaba.otter.canal.instance.manager.model.CanalParameter.*;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.DataSourcing;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.HAMode;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.IndexMode;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.MetaMode;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.SourcingType;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.StorageMode;
+import com.alibaba.otter.canal.instance.manager.model.CanalParameter.StorageScavengeMode;
+import com.alibaba.otter.canal.meta.FileMixedMetaManager;
import com.alibaba.otter.canal.meta.MemoryMetaManager;
import com.alibaba.otter.canal.meta.PeriodMixedMetaManager;
import com.alibaba.otter.canal.meta.ZooKeeperMetaManager;
@@ -32,7 +38,12 @@ import com.alibaba.otter.canal.parse.inbound.AbstractEventParser;
import com.alibaba.otter.canal.parse.inbound.group.GroupEventParser;
import com.alibaba.otter.canal.parse.inbound.mysql.LocalBinlogEventParser;
import com.alibaba.otter.canal.parse.inbound.mysql.MysqlEventParser;
-import com.alibaba.otter.canal.parse.index.*;
+import com.alibaba.otter.canal.parse.index.CanalLogPositionManager;
+import com.alibaba.otter.canal.parse.index.FailbackLogPositionManager;
+import com.alibaba.otter.canal.parse.index.MemoryLogPositionManager;
+import com.alibaba.otter.canal.parse.index.MetaLogPositionManager;
+import com.alibaba.otter.canal.parse.index.PeriodMixedLogPositionManager;
+import com.alibaba.otter.canal.parse.index.ZooKeeperLogPositionManager;
import com.alibaba.otter.canal.parse.support.AuthenticationInfo;
import com.alibaba.otter.canal.protocol.position.EntryPosition;
import com.alibaba.otter.canal.sink.entry.EntryEventSink;
@@ -110,7 +121,7 @@ public class CanalInstanceWithManager extends AbstractCanalInstance {
ZooKeeperMetaManager zooKeeperMetaManager = new ZooKeeperMetaManager();
zooKeeperMetaManager.setZkClientx(getZkclientx());
((PeriodMixedMetaManager) metaManager).setZooKeeperMetaManager(zooKeeperMetaManager);
- } else if (mode.isLocalFile()){
+ } else if (mode.isLocalFile()) {
FileMixedMetaManager fileMixedMetaManager = new FileMixedMetaManager();
fileMixedMetaManager.setDataDir(parameters.getDataDir());
fileMixedMetaManager.setPeriod(parameters.getMetaFileFlushPeriod());
diff --git a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/model/CanalParameter.java b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/model/CanalParameter.java
index befc8e2f..40df1485 100644
--- a/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/model/CanalParameter.java
+++ b/instance/manager/src/main/java/com/alibaba/otter/canal/instance/manager/model/CanalParameter.java
@@ -93,6 +93,10 @@ public class CanalParameter implements Serializable {
private Boolean filterTableError = Boolean.FALSE; // 是否忽略表解析异常
private String blackFilter = null; // 匹配黑名单,忽略解析
+ private Boolean tsdbEnable = Boolean.FALSE; // 是否开启tableMetaTSDB
+ private String tsdbJdbcUrl;
+ private String tsdbJdbcUserName;
+ private String tsdbJdbcPassword;
// ================================== 兼容字段处理
private InetSocketAddress masterAddress; // 主库信息
private String masterUsername; // 帐号
@@ -246,7 +250,7 @@ public class CanalParameter implements Serializable {
ZOOKEEPER,
/** 混合模式,内存+文件 */
MIXED,
- /** 本地文件存储模式*/
+ /** 本地文件存储模式 */
LOCAL_FILE;
public boolean isMemory() {
@@ -261,7 +265,7 @@ public class CanalParameter implements Serializable {
return this.equals(MetaMode.MIXED);
}
- public boolean isLocalFile(){
+ public boolean isLocalFile() {
return this.equals(MetaMode.LOCAL_FILE);
}
}
@@ -883,6 +887,38 @@ public class CanalParameter implements Serializable {
this.blackFilter = blackFilter;
}
+ public Boolean getTsdbEnable() {
+ return tsdbEnable;
+ }
+
+ public void setTsdbEnable(Boolean tsdbEnable) {
+ this.tsdbEnable = tsdbEnable;
+ }
+
+ public String getTsdbJdbcUrl() {
+ return tsdbJdbcUrl;
+ }
+
+ public void setTsdbJdbcUrl(String tsdbJdbcUrl) {
+ this.tsdbJdbcUrl = tsdbJdbcUrl;
+ }
+
+ public String getTsdbJdbcUserName() {
+ return tsdbJdbcUserName;
+ }
+
+ public void setTsdbJdbcUserName(String tsdbJdbcUserName) {
+ this.tsdbJdbcUserName = tsdbJdbcUserName;
+ }
+
+ public String getTsdbJdbcPassword() {
+ return tsdbJdbcPassword;
+ }
+
+ public void setTsdbJdbcPassword(String tsdbJdbcPassword) {
+ this.tsdbJdbcPassword = tsdbJdbcPassword;
+ }
+
public String toString() {
return ToStringBuilder.reflectionToString(this, CanalToStringStyle.DEFAULT_STYLE);
}
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java
index 156d4692..304ec8d7 100644
--- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/AbstractMysqlEventParser.java
@@ -8,34 +8,38 @@ import org.slf4j.LoggerFactory;
import com.alibaba.otter.canal.filter.CanalEventFilter;
import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter;
+import com.alibaba.otter.canal.parse.CanalEventParser;
import com.alibaba.otter.canal.parse.driver.mysql.packets.MysqlGTIDSet;
import com.alibaba.otter.canal.parse.exception.CanalParseException;
import com.alibaba.otter.canal.parse.inbound.AbstractEventParser;
import com.alibaba.otter.canal.parse.inbound.BinlogParser;
import com.alibaba.otter.canal.parse.inbound.MultiStageCoprocessor;
import com.alibaba.otter.canal.parse.inbound.mysql.dbsync.LogEventConvert;
+import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.DefaultTableMetaTSDBFactory;
import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.TableMetaTSDB;
-import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.TableMetaTSDBBuilder;
+import com.alibaba.otter.canal.parse.inbound.mysql.tsdb.TableMetaTSDBFactory;
import com.alibaba.otter.canal.protocol.position.EntryPosition;
public abstract class AbstractMysqlEventParser extends AbstractEventParser {
- protected final Logger logger = LoggerFactory.getLogger(this.getClass());
- protected static final long BINLOG_START_OFFEST = 4L;
+ protected final Logger logger = LoggerFactory.getLogger(this.getClass());
+ protected static final long BINLOG_START_OFFEST = 4L;
+
+ protected TableMetaTSDBFactory tableMetaTSDBFactory = new DefaultTableMetaTSDBFactory();
+ protected boolean enableTsdb = false;
+ protected String tsdbSpringXml;
+ protected TableMetaTSDB tableMetaTSDB;
- protected boolean enableTsdb = false;
- protected String tsdbSpringXml;
- protected TableMetaTSDB tableMetaTSDB;
// 编码信息
- protected byte connectionCharsetNumber = (byte) 33;
- protected Charset connectionCharset = Charset.forName("UTF-8");
- protected boolean filterQueryDcl = false;
- protected boolean filterQueryDml = false;
- protected boolean filterQueryDdl = false;
- protected boolean filterRows = false;
- protected boolean filterTableError = false;
- protected boolean useDruidDdlFilter = true;
- private final AtomicLong eventsPublishBlockingTime = new AtomicLong(0L);
+ protected byte connectionCharsetNumber = (byte) 33;
+ protected Charset connectionCharset = Charset.forName("UTF-8");
+ protected boolean filterQueryDcl = false;
+ protected boolean filterQueryDml = false;
+ protected boolean filterQueryDdl = false;
+ protected boolean filterRows = false;
+ protected boolean filterTableError = false;
+ protected boolean useDruidDdlFilter = true;
+ private final AtomicLong eventsPublishBlockingTime = new AtomicLong(0L);
protected BinlogParser buildParser() {
LogEventConvert convert = new LogEventConvert();
@@ -93,8 +97,16 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
public void start() throws CanalParseException {
if (enableTsdb) {
if (tableMetaTSDB == null) {
- // 初始化
- tableMetaTSDB = TableMetaTSDBBuilder.build(destination, tsdbSpringXml);
+ synchronized (CanalEventParser.class) {
+ try {
+ // 设置当前正在加载的通道,加载spring查找文件时会用到该变量
+ System.setProperty("canal.instance.destination", destination);
+ // 初始化
+ tableMetaTSDB = tableMetaTSDBFactory.build(destination, tsdbSpringXml);
+ } finally {
+ System.setProperty("canal.instance.destination", "");
+ }
+ }
}
}
@@ -103,7 +115,7 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
public void stop() throws CanalParseException {
if (enableTsdb) {
- TableMetaTSDBBuilder.destory(destination);
+ tableMetaTSDBFactory.destory(destination);
tableMetaTSDB = null;
}
@@ -177,7 +189,7 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
if (this.enableTsdb) {
if (tableMetaTSDB == null) {
// 初始化
- tableMetaTSDB = TableMetaTSDBBuilder.build(destination, tsdbSpringXml);
+ tableMetaTSDB = tableMetaTSDBFactory.build(destination, tsdbSpringXml);
}
}
}
@@ -187,7 +199,7 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
if (this.enableTsdb) {
if (tableMetaTSDB == null) {
// 初始化
- tableMetaTSDB = TableMetaTSDBBuilder.build(destination, tsdbSpringXml);
+ tableMetaTSDB = tableMetaTSDBFactory.build(destination, tsdbSpringXml);
}
}
}
@@ -195,5 +207,9 @@ public abstract class AbstractMysqlEventParser extends AbstractEventParser {
public AtomicLong getEventsPublishBlockingTime() {
return this.eventsPublishBlockingTime;
}
+
+ public void setTableMetaTSDBFactory(TableMetaTSDBFactory tableMetaTSDBFactory) {
+ this.tableMetaTSDBFactory = tableMetaTSDBFactory;
+ }
}
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java
index 69cf0572..7db700ae 100644
--- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/MysqlEventParser.java
@@ -63,7 +63,6 @@ public class MysqlEventParser extends AbstractMysqlEventParser implements CanalE
private String detectingSQL; // 心跳sql
private MysqlConnection metaConnection; // 查询meta信息的链接
private TableMetaCache tableMetaCache; // 对应meta
- // cache
private int fallbackIntervalInSeconds = 60; // 切换回退时间
private BinlogFormat[] supportBinlogFormats; // 支持的binlogFormat,如果设置会执行强校验
private BinlogImage[] supportBinlogImages; // 支持的binlogImage,如果设置会执行强校验
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java
index 09af808f..dc8e57d7 100644
--- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/dbsync/LogEventConvert.java
@@ -541,12 +541,16 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
tableError |= parseOneRow(rowDataBuilder, event, buffer, changeColumns, true, tableMeta);
}
- rowsCount ++;
+ rowsCount++;
rowChangeBuider.addRowDatas(rowDataBuilder.build());
}
TableMapLogEvent table = event.getTable();
- Header header = createHeader(event.getHeader(), table.getDbName(), table.getTableName(), eventType, rowsCount);
+ Header header = createHeader(event.getHeader(),
+ table.getDbName(),
+ table.getTableName(),
+ eventType,
+ rowsCount);
RowChange rowChange = rowChangeBuider.build();
if (tableError) {
@@ -801,12 +805,12 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
return createEntry(header, EntryType.ROWDATA, rowChangeBuider.build().toByteString());
}
-
private Header createHeader(LogHeader logHeader, String schemaName, String tableName, EventType eventType) {
return createHeader(logHeader, schemaName, tableName, eventType, -1);
}
- private Header createHeader(LogHeader logHeader, String schemaName, String tableName, EventType eventType, Integer rowsCount) {
+ private Header createHeader(LogHeader logHeader, String schemaName, String tableName, EventType eventType,
+ Integer rowsCount) {
// header会做信息冗余,方便以后做检索或者过滤
Header.Builder headerBuilder = Header.newBuilder();
headerBuilder.setVersion(version);
@@ -960,5 +964,4 @@ public class LogEventConvert extends AbstractCanalLifeCycle implements BinlogPar
public void setGtidSet(GTIDSet gtidSet) {
this.gtidSet = gtidSet;
}
-
}
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DefaultTableMetaTSDBFactory.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DefaultTableMetaTSDBFactory.java
new file mode 100644
index 00000000..00e2744b
--- /dev/null
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/DefaultTableMetaTSDBFactory.java
@@ -0,0 +1,19 @@
+package com.alibaba.otter.canal.parse.inbound.mysql.tsdb;
+
+/**
+ * @author agapple 2017年10月11日 下午8:45:40
+ * @since 1.0.25
+ */
+public class DefaultTableMetaTSDBFactory implements TableMetaTSDBFactory {
+
+ /**
+ * 代理一下tableMetaTSDB的获取,使用隔离的spring定义
+ */
+ public TableMetaTSDB build(String destination, String springXml) {
+ return TableMetaTSDBBuilder.build(destination, springXml);
+ }
+
+ public void destory(String destination) {
+ TableMetaTSDBBuilder.destory(destination);
+ }
+}
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBBuilder.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBBuilder.java
index 8e37ef71..107e1013 100644
--- a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBBuilder.java
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBBuilder.java
@@ -10,12 +10,15 @@ import org.springframework.context.support.ClassPathXmlApplicationContext;
import com.google.common.collect.Maps;
/**
- * @author agapple 2017年10月11日 下午8:45:40
+ * tableMeta构造器
+ *
+ * @author agapple 2018年8月8日 上午11:01:08
* @since 1.0.25
*/
+
public class TableMetaTSDBBuilder {
- protected final static Logger logger = LoggerFactory.getLogger(TableMetaTSDBBuilder.class);
+ protected final static Logger logger = LoggerFactory.getLogger(DefaultTableMetaTSDBFactory.class);
private static ConcurrentMap contexts = Maps.newConcurrentMap();
/**
diff --git a/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBFactory.java b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBFactory.java
new file mode 100644
index 00000000..950645a8
--- /dev/null
+++ b/parse/src/main/java/com/alibaba/otter/canal/parse/inbound/mysql/tsdb/TableMetaTSDBFactory.java
@@ -0,0 +1,18 @@
+package com.alibaba.otter.canal.parse.inbound.mysql.tsdb;
+
+/**
+ * tableMeta构造器,允许重载实现
+ *
+ * @author agapple 2018年8月8日 上午11:01:08
+ * @since 1.0.26
+ */
+
+public interface TableMetaTSDBFactory {
+
+ /**
+ * 代理一下tableMetaTSDB的获取,使用隔离的spring定义
+ */
+ public TableMetaTSDB build(String destination, String springXml);
+
+ public void destory(String destination);
+}
diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/NoStorageTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/NoStorageTest.java
new file mode 100644
index 00000000..04fd95ea
--- /dev/null
+++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/NoStorageTest.java
@@ -0,0 +1,38 @@
+package com.alibaba.otter.canal.parse.inbound.mysql.tablemeta;
+
+import com.alibaba.otter.canal.parse.inbound.TableMeta;
+import com.alibaba.otter.canal.parse.inbound.mysql.MysqlConnection;
+import com.alibaba.otter.canal.protocol.position.EntryPosition;
+import org.junit.Test;
+
+import java.net.InetSocketAddress;
+import java.util.Date;
+
+public class NoStorageTest {
+ final String DBNAME = "testdb";
+ final String TBNAME = "testtb";
+ final String DDL = "CREATE TABLE `testtb` (\n" +
+ " `id` int(11) NOT NULL AUTO_INCREMENT,\n" +
+ " `name` varchar(2048) DEFAULT NULL,\n" +
+ " `datachange_lasttime` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '最晚更新时间',\n" +
+ " `otter_testcol` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol1` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol2` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol3` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol4` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol5` varchar(45) DEFAULT NULL,\n" +
+ " PRIMARY KEY (`id`)\n" +
+ " ) ENGINE=InnoDB AUTO_INCREMENT=58333898 DEFAULT CHARSET=utf8mb4";
+ @Test
+ public void nostorage() {
+ MysqlConnection connection = new MysqlConnection(new InetSocketAddress("127.0.0.1", 3306), "root", "hello");
+ TableMetaCacheWithStorage tableMetaCacheWithStorage = new TableMetaCacheWithStorage(connection, null);
+ EntryPosition entryPosition = new EntryPosition();
+ entryPosition.setTimestamp(new Date().getTime());
+ String fullTableName = DBNAME + "." + TBNAME;
+ tableMetaCacheWithStorage.apply(entryPosition, fullTableName, DDL, null);
+ entryPosition.setTimestamp(new Date().getTime() + 1000L);
+ TableMeta result = tableMetaCacheWithStorage.getTableMeta(DBNAME, TBNAME, false, entryPosition);
+ assert result.getDdl().equalsIgnoreCase(DDL);
+ }
+}
diff --git a/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/StorageTest.java b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/StorageTest.java
new file mode 100644
index 00000000..0ec16607
--- /dev/null
+++ b/parse/src/test/java/com/alibaba/otter/canal/parse/inbound/mysql/tablemeta/StorageTest.java
@@ -0,0 +1,77 @@
+package com.alibaba.otter.canal.parse.inbound.mysql.tablemeta;
+
+import com.alibaba.otter.canal.parse.inbound.TableMeta;
+import com.alibaba.otter.canal.parse.inbound.mysql.MysqlConnection;
+import com.alibaba.otter.canal.parse.inbound.mysql.tablemeta.impl.mysql.MySqlTableMetaCallback;
+import com.alibaba.otter.canal.parse.inbound.mysql.tablemeta.impl.mysql.MySqlTableMetaStorageFactory;
+import com.alibaba.otter.canal.protocol.position.EntryPosition;
+import com.alibaba.otter.canal.protocol.position.Position;
+import org.junit.Test;
+
+import java.net.InetSocketAddress;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+
+public class StorageTest {
+
+ final String DBNAME = "testdb";
+ final String TBNAME = "testtb";
+ final String DDL = "CREATE TABLE `testtb` (\n" +
+ " `id` int(11) NOT NULL AUTO_INCREMENT,\n" +
+ " `name` varchar(2048) DEFAULT NULL,\n" +
+ " `datachange_lasttime` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '最晚更新时间',\n" +
+ " `otter_testcol` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol1` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol2` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol3` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol4` varchar(45) DEFAULT NULL,\n" +
+ " `otter_testcol5` varchar(45) DEFAULT NULL,\n" +
+ " PRIMARY KEY (`id`)\n" +
+ " ) ENGINE=InnoDB AUTO_INCREMENT=58333898 DEFAULT CHARSET=utf8mb4";
+
+ @Test
+ public void storage() {
+
+ MySqlTableMetaStorageFactory factory = new MySqlTableMetaStorageFactory(new MySqlTableMetaCallback() {
+ @Override
+ public void save(String dbAddress, String schema, String table, String ddl, Long timestamp) {
+
+ }
+
+ @Override
+ public List fetch(String dbAddress, String dbName) {
+ TableMetaEntry tableMeta = new TableMetaEntry();
+ tableMeta.setSchema(DBNAME);
+ tableMeta.setTable(TBNAME);
+ tableMeta.setDdl(DDL);
+ tableMeta.setTimestamp(new Date().getTime());
+ List entries = new ArrayList();
+ entries.add(tableMeta);
+ return entries;
+ }
+
+ @Override
+ public List fetch(String dbAddress, String dbName, String tableName) {
+ TableMetaEntry tableMeta = new TableMetaEntry();
+ tableMeta.setSchema(DBNAME);
+ tableMeta.setTable(TBNAME);
+ tableMeta.setDdl(DDL);
+ tableMeta.setTimestamp(new Date().getTime());
+ List entries = new ArrayList();
+ entries.add(tableMeta);
+ return entries;
+ }
+ }, DBNAME);
+ MysqlConnection connection = new MysqlConnection(new InetSocketAddress("127.0.0.1", 3306), "root", "hello");
+ TableMetaCacheWithStorage tableMetaCacheWithStorage = new TableMetaCacheWithStorage(connection, factory.getTableMetaStorage());
+ EntryPosition entryPosition = new EntryPosition();
+ entryPosition.setTimestamp(new Date().getTime());
+ String fullTableName = DBNAME + "." + TBNAME;
+ tableMetaCacheWithStorage.apply(entryPosition, fullTableName, DDL, null);
+
+ entryPosition.setTimestamp(new Date().getTime() + 1000L);
+ TableMeta result = tableMetaCacheWithStorage.getTableMeta(DBNAME, TBNAME, false, entryPosition);
+ assert result.getDdl().equalsIgnoreCase(DDL);
+ }
+}
diff --git a/pom.xml b/pom.xml
index 9a14868b..56ac821a 100644
--- a/pom.xml
+++ b/pom.xml
@@ -96,8 +96,8 @@
true
true
- 1.6
- 1.6
+ 1.7
+ 1.7
UTF-8
3.2.9.RELEASE
@@ -247,7 +247,7 @@
com.alibaba.fastsql
fastsql
- 2.0.0_preview_520
+ 2.0.0_preview_540
com.alibaba
@@ -332,7 +332,7 @@
org.apache.maven.plugins
maven-compiler-plugin
- 3.2
+ 3.8.0
${java_source_version}
${java_target_version}