如果Transformer/Merger/Loader线程执行较慢,那Kafka Offset就迟迟不能提交。
这种情况下,应该会造成重复消费到同样的数据吧。具体见如下代码:
|
public ChangeSet pollChangeSet() throws BiremeException { |
|
ConsumerRecords<String, String> records = null; |
|
|
|
try { |
|
records = consumer.poll(POLL_TIMEOUT); |
|
} catch (InterruptException e) { |
|
} |
|
|
Ps: 看来bireme只能依赖主键做兜底了。
如果Transformer/Merger/Loader线程执行较慢,那Kafka Offset就迟迟不能提交。
这种情况下,应该会造成重复消费到同样的数据吧。具体见如下代码:
bireme/src/main/java/cn/hashdata/bireme/pipeline/KafkaPipeLine.java
Lines 44 to 51 in 9cfc128
Ps: 看来bireme只能依赖主键做兜底了。