Я настраиваю процессор потока Flink, используя kafka иasticsearch. Я хочу воспроизвести свои данные, но когда я устанавливаю параллелизм более чем на 1, это не завершает sh программу, которую я считаю такой, потому что поток 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 об окончании после завершения одного потока?