@Override
protected int getRecordsInTarget() {
int expectedRecordsInTarget = 0;
for(KafkaStream<byte[], byte[]> kafkaStream : kafkaStreams) {
ConsumerIterator<byte[], byte[]> it = kafkaStream.iterator();
try {
while (it.hasNext()) {
expectedRecordsInTarget++;
it.next();
}
} catch (kafka.consumer.ConsumerTimeoutException e) {
//no-op
}
}
return expectedRecordsInTarget;
}
KafkaDestinationMultiPartitionPipelineRunIT.java 文件源码
java
阅读 15
收藏 0
点赞 0
评论 0
项目:datacollector
作者:
评论列表
文章目录