Merge branch 'master' into feature/mcy

This commit is contained in:
mcy
2018-11-01 14:53:51 +08:00
@@ -77,8 +77,8 @@ public class CanalRocketMQProducer implements CanalMQProducer {
JSON.toJSONString(flatMessage),
destination.getTopic(),
destination.getPartition());
Message message = new Message(destination.getTopic(), JSON.toJSONString(flatMessage)
.getBytes());
Message message = new Message(destination.getTopic(),
JSON.toJSONString(flatMessage).getBytes());
this.defaultMQProducer.send(message, new MessageQueueSelector() {
@Override
@@ -98,28 +98,32 @@ public class CanalRocketMQProducer implements CanalMQProducer {
int length = partitionFlatMessage.length;
for (int i = 0; i < length; i++) {
FlatMessage flatMessagePart = partitionFlatMessage[i];
logger.debug("flatMessagePart: {}, partition: {}",
JSON.toJSONString(flatMessagePart),
i);
final int index = i;
try {
Message message = new Message(destination.getTopic(),
JSON.toJSONString(flatMessagePart).getBytes());
this.defaultMQProducer.send(message, new MessageQueueSelector() {
if (flatMessagePart != null) {
logger.debug("flatMessagePart: {}, partition: {}",
JSON.toJSONString(flatMessagePart),
i);
final int index = i;
try {
Message message = new Message(destination.getTopic(),
JSON.toJSONString(flatMessagePart).getBytes());
this.defaultMQProducer.send(message, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
if (index > mqs.size()) {
throw new CanalServerException("partition number is error,config num:"
+ destination.getPartitionsNum()
+ ", mq num: " + mqs.size());
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg,
Object arg) {
if (index > mqs.size()) {
throw new CanalServerException(
"partition number is error,config num:"
+ destination.getPartitionsNum()
+ ", mq num: " + mqs.size());
}
return mqs.get(index);
}
return mqs.get(index);
}
}, null);
} catch (Exception e) {
logger.error("send flat message to hashed partition error", e);
callback.rollback();
}, null);
} catch (Exception e) {
logger.error("send flat message to hashed partition error", e);
callback.rollback();
}
}
}
}