kafka 客户端 example 完善

This commit is contained in:
rewerma
2018-06-13 14:01:30 +08:00
parent ec072a45d8
commit 9964bb75b1
@@ -86,40 +86,42 @@ public class KafkaCanalClientExample {
}
private void process() {
while (running) {
try {
connector.subscribe();
while (running) {
try {
Message message = connector.getWithoutAck(1L, TimeUnit.SECONDS); //获取message
if (message == null) {
continue;
}
long batchId = message.getId();
int size = message.getEntries().size();
if (batchId == -1 || size == 0) {
// try {
// Thread.sleep(1000);
// } catch (InterruptedException e) {
// }
} else {
// printSummary(message, batchId, size);
// printEntry(message.getEntries());
logger.info(message.toString());
}
connector.ack(); // 提交确认
} catch (WakeupException e) {
//ignore
} catch (Exception e) {
logger.error(e.getMessage(), e);
try {
connector.subscribe();
while (running) {
try {
Message message = connector.getWithoutAck(1L, TimeUnit.SECONDS); //获取message
if (message == null) {
continue;
}
long batchId = message.getId();
int size = message.getEntries().size();
if (batchId == -1 || size == 0) {
// try {
// Thread.sleep(1000);
// } catch (InterruptedException e) {
// }
} else {
// printSummary(message, batchId, size);
// printEntry(message.getEntries());
logger.info(message.toString());
}
connector.ack(); // 提交确认
} catch (WakeupException e) {
try {
Thread.sleep(500); //延时确保running状态的改变
} catch (InterruptedException ie) {
//ignore
}
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
} catch (WakeupException e) {
//ignore
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
} catch (WakeupException e) {
//ignore
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
}
}