【问题标题】:Flink stream not finishingFlink 流未完成
【发布时间】: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 消费者组上的所有线程在一个线程完成后结束?

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    如果您实现了自己的 SourceFunction,请使用 cancel 方法,如 Flink SourceFunction 中的示例所示。 FlinkKafkaConsumerBase 类也有 cancel 方法。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-06-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-30
      • 2018-09-25
      • 2012-07-28
      • 2019-01-26
      相关资源
      最近更新 更多