【发布时间】:2020-06-28 14:15:12
【问题描述】:
我正在使用 kafka 和 elasticsearch 设置一个 flink 流处理器。我想重播我的数据,但是当我将并行度设置为大于 1 时,它没有完成程序我相信这是因为 kafka 流只看到一条消息被识别为流的结尾。
public CustomSchema(Date _endTime) {
endTime = _endTime;
}
@Override
public boolean isEndOfStream(CustomTopicWrapper nextElement) {
if (this.endTime != null && nextElement.messageTime.getTime() >= this.endTime.getTime()) {
return true;
}
return false;
}
有没有办法告诉 flink 消费者组上的所有线程在一个线程完成后结束?
【问题讨论】: