>();
- 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/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java b/example/src/main/java/com/alibaba/otter/canal/example/kafka/AbstractKafkaTest.java
similarity index 74%
rename from client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java
rename to example/src/main/java/com/alibaba/otter/canal/example/kafka/AbstractKafkaTest.java
index 473baf67..b73db92e 100644
--- a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/AbstractKafkaTest.java
+++ b/example/src/main/java/com/alibaba/otter/canal/example/kafka/AbstractKafkaTest.java
@@ -1,6 +1,6 @@
-package com.alibaba.otter.canal.client.running.kafka;
+package com.alibaba.otter.canal.example.kafka;
-import org.junit.Assert;
+import com.alibaba.otter.canal.example.BaseCanalClientTest;
/**
* Kafka 测试基类
@@ -8,7 +8,7 @@ import org.junit.Assert;
* @author machengyuan @ 2018-6-12
* @version 1.0.0
*/
-public abstract class AbstractKafkaTest {
+public abstract class AbstractKafkaTest extends BaseCanalClientTest {
public static String topic = "example";
public static Integer partition = null;
@@ -20,7 +20,6 @@ public abstract class AbstractKafkaTest {
try {
Thread.sleep(time);
} catch (InterruptedException e) {
- Assert.fail(e.getMessage());
}
}
}
diff --git a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientExample.java
similarity index 93%
rename from client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java
rename to example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientExample.java
index bf739a06..83757aec 100644
--- a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaClientExample.java
+++ b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaClientExample.java
@@ -1,9 +1,8 @@
-package com.alibaba.otter.canal.client.running.kafka;
+package com.alibaba.otter.canal.example.kafka;
import java.util.List;
import java.util.concurrent.TimeUnit;
-import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.Assert;
@@ -40,8 +39,7 @@ public class CanalKafkaClientExample {
public static void main(String[] args) {
try {
- final CanalKafkaClientExample kafkaCanalClientExample = new CanalKafkaClientExample(
- AbstractKafkaTest.zkServers,
+ final CanalKafkaClientExample kafkaCanalClientExample = new CanalKafkaClientExample(AbstractKafkaTest.zkServers,
AbstractKafkaTest.servers,
AbstractKafkaTest.topic,
AbstractKafkaTest.partition,
@@ -141,11 +139,7 @@ public class CanalKafkaClientExample {
}
}
- try {
- connector.unsubscribe();
- } catch (WakeupException e) {
- // No-op. Continue process
- }
+ connector.unsubscribe();
connector.disconnect();
}
}
diff --git a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaOffsetClientExample.java b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaOffsetClientExample.java
similarity index 79%
rename from client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaOffsetClientExample.java
rename to example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaOffsetClientExample.java
index ff45e88f..22f6e772 100644
--- a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/CanalKafkaOffsetClientExample.java
+++ b/example/src/main/java/com/alibaba/otter/canal/example/kafka/CanalKafkaOffsetClientExample.java
@@ -1,9 +1,8 @@
-package com.alibaba.otter.canal.client.running.kafka;
+package com.alibaba.otter.canal.example.kafka;
import java.util.List;
import java.util.concurrent.TimeUnit;
-import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.Assert;
@@ -13,40 +12,46 @@ import com.alibaba.otter.canal.client.kafka.protocol.KafkaMessage;
/**
* KafkaOffsetCanalConnector 使用示例
- *
KafkaOffsetCanalConnector 与 KafkaCanalConnector 的另一区别是 auto.offset.reset 默认值不同;
- * KafkaOffsetCanalConnector 默认为 earliest;canal-kafka-client重启后从未被消费的记录开始拉取消息,同时提供了修改 auto.offset.reset 的方法 setAutoOffsetReset
+ *
+ * KafkaOffsetCanalConnector 与 KafkaCanalConnector 的另一区别是 auto.offset.reset
+ * 默认值不同;
+ *
+ *
+ * KafkaOffsetCanalConnector 默认为
+ * earliest;canal-kafka-client重启后从未被消费的记录开始拉取消息,同时提供了修改 auto.offset.reset 的方法
+ * setAutoOffsetReset
+ *
*
* @author panjianping @ 2018-12-18
* @version 1.1.3
*/
public class CanalKafkaOffsetClientExample {
- protected final static Logger logger = LoggerFactory.getLogger(CanalKafkaOffsetClientExample.class);
+ protected final static Logger logger = LoggerFactory.getLogger(CanalKafkaOffsetClientExample.class);
- private KafkaOffsetCanalConnector connector;
+ private KafkaOffsetCanalConnector connector;
- private static volatile boolean running = false;
+ private static volatile boolean running = false;
- private Thread thread = null;
+ private Thread thread = null;
private Thread.UncaughtExceptionHandler handler = new Thread.UncaughtExceptionHandler() {
- public void uncaughtException(Thread t, Throwable e) {
- logger.error("parse events has an error", e);
- }
- };
+ public void uncaughtException(Thread t, Throwable e) {
+ logger.error("parse events has an error", e);
+ }
+ };
- public CanalKafkaOffsetClientExample(String servers, String topic, Integer partition, String groupId) {
+ public CanalKafkaOffsetClientExample(String servers, String topic, Integer partition, String groupId){
connector = new KafkaOffsetCanalConnector(servers, topic, partition, groupId, false);
}
public static void main(String[] args) {
try {
- final CanalKafkaOffsetClientExample kafkaCanalClientExample = new CanalKafkaOffsetClientExample(
- AbstractKafkaTest.servers,
- AbstractKafkaTest.topic,
- AbstractKafkaTest.partition,
- AbstractKafkaTest.groupId);
+ final CanalKafkaOffsetClientExample kafkaCanalClientExample = new CanalKafkaOffsetClientExample(AbstractKafkaTest.servers,
+ AbstractKafkaTest.topic,
+ AbstractKafkaTest.partition,
+ AbstractKafkaTest.groupId);
logger.info("## start the kafka consumer: {}-{}", AbstractKafkaTest.topic, AbstractKafkaTest.groupId);
kafkaCanalClientExample.start();
logger.info("## the canal kafka consumer is running now ......");
@@ -110,7 +115,7 @@ public class CanalKafkaOffsetClientExample {
while (running) {
try {
// 修改 AutoOffsetReset 的值,默认(earliest)
- //connector.setAutoOffsetReset(null);
+ // connector.setAutoOffsetReset(null);
connector.connect();
connector.subscribe();
// 消息起始偏移地址
@@ -164,11 +169,7 @@ public class CanalKafkaOffsetClientExample {
}
}
- try {
- connector.unsubscribe();
- } catch (WakeupException e) {
- // No-op. Continue process
- }
+ connector.unsubscribe();
connector.disconnect();
}
}
diff --git a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java b/example/src/main/java/com/alibaba/otter/canal/example/kafka/KafkaClientRunningTest.java
similarity index 72%
rename from client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java
rename to example/src/main/java/com/alibaba/otter/canal/example/kafka/KafkaClientRunningTest.java
index 9133f74c..f1e12cd8 100644
--- a/client/src/test/java/com/alibaba/otter/canal/client/running/kafka/KafkaClientRunningTest.java
+++ b/example/src/main/java/com/alibaba/otter/canal/example/kafka/KafkaClientRunningTest.java
@@ -1,12 +1,10 @@
-package com.alibaba.otter.canal.client.running.kafka;
+package com.alibaba.otter.canal.example.kafka;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
-import org.apache.kafka.common.errors.WakeupException;
-import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -25,12 +23,9 @@ public class KafkaClientRunningTest extends AbstractKafkaTest {
private boolean running = true;
- @Test
public void testKafkaConsumer() {
final ExecutorService executor = Executors.newFixedThreadPool(1);
-
final KafkaCanalConnector connector = new KafkaCanalConnector(servers, topic, partition, groupId, null, false);
-
executor.submit(new Runnable() {
@Override
@@ -38,15 +33,11 @@ public class KafkaClientRunningTest extends AbstractKafkaTest {
connector.connect();
connector.subscribe();
while (running) {
- try {
- List messages = connector.getList(3L, TimeUnit.SECONDS);
- if (messages != null) {
- System.out.println(messages);
- }
- connector.ack();
- } catch (WakeupException e) {
- // ignore
+ List messages = connector.getList(3L, TimeUnit.SECONDS);
+ if (messages != null) {
+ System.out.println(messages);
}
+ connector.ack();
}
connector.unsubscribe();
connector.disconnect();
diff --git a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java
new file mode 100644
index 00000000..55672db1
--- /dev/null
+++ b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/AbstractRocektMQTest.java
@@ -0,0 +1,11 @@
+package com.alibaba.otter.canal.example.rocketmq;
+
+import com.alibaba.otter.canal.example.BaseCanalClientTest;
+
+public abstract class AbstractRocektMQTest extends BaseCanalClientTest {
+
+ public static String topic = "example";
+ public static String groupId = "group";
+ public static String nameServers = "localhost:9876";
+
+}
diff --git a/client/src/test/java/com/alibaba/otter/canal/client/running/rocketmq/CanalRocketMQClientExample.java b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java
similarity index 88%
rename from client/src/test/java/com/alibaba/otter/canal/client/running/rocketmq/CanalRocketMQClientExample.java
rename to example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java
index 2eaa398c..c95abb01 100644
--- a/client/src/test/java/com/alibaba/otter/canal/client/running/rocketmq/CanalRocketMQClientExample.java
+++ b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientExample.java
@@ -1,15 +1,14 @@
-package com.alibaba.otter.canal.client.running.rocketmq;
+package com.alibaba.otter.canal.example.rocketmq;
import java.util.List;
import java.util.concurrent.TimeUnit;
-import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.util.Assert;
import com.alibaba.otter.canal.client.rocketmq.RocketMQCanalConnector;
-import com.alibaba.otter.canal.client.running.kafka.AbstractKafkaTest;
+import com.alibaba.otter.canal.example.kafka.AbstractKafkaTest;
import com.alibaba.otter.canal.protocol.Message;
/**
@@ -119,9 +118,9 @@ public class CanalRocketMQClientExample extends AbstractRocektMQTest {
// } catch (InterruptedException e) {
// }
} else {
- // printSummary(message, batchId, size);
- // printEntry(message.getEntries());
- logger.info(message.toString());
+ printSummary(message, batchId, size);
+ printEntry(message.getEntries());
+ // logger.info(message.toString());
}
}
@@ -132,11 +131,7 @@ public class CanalRocketMQClientExample extends AbstractRocektMQTest {
}
}
- try {
- connector.unsubscribe();
- } catch (WakeupException e) {
- // No-op. Continue process
- }
-// connector.stopRunning();
+ connector.unsubscribe();
+ // connector.stopRunning();
}
}
diff --git a/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientFlatMessageExample.java b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientFlatMessageExample.java
new file mode 100644
index 00000000..6b553771
--- /dev/null
+++ b/example/src/main/java/com/alibaba/otter/canal/example/rocketmq/CanalRocketMQClientFlatMessageExample.java
@@ -0,0 +1,134 @@
+package com.alibaba.otter.canal.example.rocketmq;
+
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.util.Assert;
+
+import com.alibaba.otter.canal.client.rocketmq.RocketMQCanalConnector;
+import com.alibaba.otter.canal.example.kafka.AbstractKafkaTest;
+import com.alibaba.otter.canal.protocol.FlatMessage;
+
+/**
+ * Kafka client example
+ *
+ * @author machengyuan @ 2018-6-12
+ * @version 1.0.0
+ */
+public class CanalRocketMQClientFlatMessageExample extends AbstractRocektMQTest {
+
+ protected final static Logger logger = LoggerFactory.getLogger(CanalRocketMQClientFlatMessageExample.class);
+
+ private RocketMQCanalConnector connector;
+
+ private static volatile boolean running = false;
+
+ private Thread thread = null;
+
+ private Thread.UncaughtExceptionHandler handler = new Thread.UncaughtExceptionHandler() {
+
+ public void uncaughtException(Thread t, Throwable e) {
+ logger.error("parse events has an error", e);
+ }
+ };
+
+ public CanalRocketMQClientFlatMessageExample(String nameServers, String topic, String groupId){
+ connector = new RocketMQCanalConnector(nameServers, topic, groupId, 500, true);
+ }
+
+ public static void main(String[] args) {
+ try {
+ final CanalRocketMQClientFlatMessageExample rocketMQClientExample = new CanalRocketMQClientFlatMessageExample(nameServers,
+ topic,
+ groupId);
+ logger.info("## Start the rocketmq consumer: {}-{}", AbstractKafkaTest.topic, AbstractKafkaTest.groupId);
+ rocketMQClientExample.start();
+ logger.info("## The canal rocketmq consumer is running now ......");
+ Runtime.getRuntime().addShutdownHook(new Thread() {
+
+ public void run() {
+ try {
+ logger.info("## Stop the rocketmq consumer");
+ rocketMQClientExample.stop();
+ } catch (Throwable e) {
+ logger.warn("## Something goes wrong when stopping rocketmq consumer:", e);
+ } finally {
+ logger.info("## Rocketmq consumer is down.");
+ }
+ }
+
+ });
+ while (running)
+ ;
+ } catch (Throwable e) {
+ logger.error("## Something going wrong when starting up the rocketmq consumer:", e);
+ System.exit(0);
+ }
+ }
+
+ public void start() {
+ Assert.notNull(connector, "connector is null");
+ thread = new Thread(new Runnable() {
+
+ public void run() {
+ process();
+ }
+ });
+ thread.setUncaughtExceptionHandler(handler);
+ thread.start();
+ running = true;
+ }
+
+ public void stop() {
+ if (!running) {
+ return;
+ }
+ running = false;
+ if (thread != null) {
+ try {
+ thread.join();
+ } catch (InterruptedException e) {
+ // ignore
+ }
+ }
+ }
+
+ private void process() {
+ while (!running) {
+ try {
+ Thread.sleep(1000);
+ } catch (InterruptedException e) {
+ }
+ }
+
+ while (running) {
+ try {
+ connector.connect();
+ connector.subscribe();
+ while (running) {
+ List messages = connector.getFlatList(100L, TimeUnit.MILLISECONDS); // 获取message
+ for (FlatMessage message : messages) {
+ long batchId = message.getId();
+ if (batchId == -1 || message.getData() == null) {
+ // try {
+ // Thread.sleep(1000);
+ // } catch (InterruptedException e) {
+ // }
+ } else {
+ logger.info(message.toString());
+ }
+ }
+
+ connector.ack(); // 提交确认
+ }
+ } catch (Exception e) {
+ logger.error(e.getMessage(), e);
+ }
+ }
+
+ connector.unsubscribe();
+ // connector.stopRunning();
+ }
+}
diff --git a/pom.xml b/pom.xml
index 266b1921..5de04ea4 100644
--- a/pom.xml
+++ b/pom.xml
@@ -294,7 +294,7 @@
mysql
mysql-connector-java
- 5.1.40
+ 5.1.47
diff --git a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java
index 163dcade..370b848d 100644
--- a/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java
+++ b/server/src/main/java/com/alibaba/otter/canal/common/MQMessageUtils.java
@@ -13,6 +13,7 @@ import org.apache.commons.lang.StringUtils;
import com.alibaba.otter.canal.filter.aviater.AviaterRegexFilter;
import com.alibaba.otter.canal.protocol.CanalEntry;
import com.alibaba.otter.canal.protocol.CanalEntry.Entry;
+import com.alibaba.otter.canal.protocol.CanalEntry.RowChange;
import com.alibaba.otter.canal.protocol.FlatMessage;
import com.alibaba.otter.canal.protocol.Message;
import com.google.common.base.Function;
@@ -207,6 +208,12 @@ public class MQMessageUtils {
if (hashMode == null) {
// 如果都没有匹配,发送到第一个分区
partitionEntries[0].add(entry);
+ } else if (hashMode.tableHash) {
+ int hashCode = table.hashCode();
+ int pkHash = Math.abs(hashCode) % partitionsNum;
+ pkHash = Math.abs(pkHash);
+ // tableHash not need split entry message
+ partitionEntries[pkHash].add(entry);
} else {
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
int hashCode = table.hashCode();
@@ -217,7 +224,7 @@ public class MQMessageUtils {
hashCode = hashCode ^ column.getValue().hashCode();
}
}
- } else if (!hashMode.tableHash) {
+ } else {
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
if (checkPkNamesHasContain(hashMode.pkNames, column.getName())) {
hashCode = hashCode ^ column.getValue().hashCode();
@@ -227,7 +234,14 @@ public class MQMessageUtils {
int pkHash = Math.abs(hashCode) % partitionsNum;
pkHash = Math.abs(pkHash);
- partitionEntries[pkHash].add(entry);
+ // build new entry
+ Entry.Builder builder = Entry.newBuilder(entry);
+ RowChange.Builder rowChangeBuilder = RowChange.newBuilder(rowChange);
+ rowChangeBuilder.clearRowDatas();
+ rowChangeBuilder.addRowDatas(rowData);
+ builder.clearStoreValue();
+ builder.setStoreValue(rowChangeBuilder.build().toByteString());
+ partitionEntries[pkHash].add(builder.build());
}
}
} else {
@@ -305,6 +319,7 @@ public class MQMessageUtils {
List