【问题标题】:Closing Spark Streaming Context after first batch (trying to retrieve kafka offsets)第一批后关闭 Spark Streaming Context(尝试检索 kafka 偏移量)
【发布时间】:2018-12-11 00:39:13
【问题描述】:

我正在尝试为我的 Spark Ba​​tch 作业检索 Kafka 偏移量。检索偏移量后,我想关闭流上下文。

我尝试向流上下文添加一个流侦听器,并实现 onBatchCompleted 方法以在作业完成后关闭流,但我收到异常“无法在侦听器总线线程中停止 StreamingContext”

有解决办法吗?我正在尝试检索偏移量以调用 KafkaUtils.createRDD(sparkContext, kafkaProperties, OffsetRange[], LocationStrateg)

private OffsetRange[] getOffsets(SparkConf sparkConf) throws InterruptedException {
    final AtomicReference<OffsetRange[]> atomicReference = new AtomicReference<>();

    JavaStreamingContext sc = new JavaStreamingContext(sparkConf, Duration.apply(50));
    JavaInputDStream<ConsumerRecord<String, String>> stream =
            KafkaUtils.createDirectStream(sc, LocationStrategies.PreferConsistent(), ConsumerStrategies.<String, String>Subscribe(Arrays.asList("test"), getKafkaParam()));
    stream.foreachRDD((VoidFunction<JavaRDD<ConsumerRecord<String, String>>>) rdd -> {
                atomicReference.set(((HasOffsetRanges) rdd.rdd()).offsetRanges());
                // sc.stop(false); //this would throw exception saying consumer is already closed
            }
    );
    sc.addStreamingListener(new TopicListener(sc)); //Throws exception saying "Cannot stop StreamingContext within listener bus thread."
    sc.start();
    sc.awaitTermination();
    return atomicReference.get();
}



public class TopicListener implements StreamingListener {
private JavaStreamingContext sc;

public TopicListener(JavaStreamingContext sc){
    this.sc = sc;
}
@Override
public void onBatchCompleted(StreamingListenerBatchCompleted streamingListenerBatchCompleted) {
    sc.stop(false);
}

非常感谢stackoverflow-ers :) 我已经尝试搜索可能的解决方案,但到目前为止还没有成功

编辑: 我使用 KafkaConsumer 来获取分区信息。获得分区信息后,我会创建一个 TopicPartition pojos 列表并调用 position 和 endOffsets 方法来分别获取我的 groupId 的当前位置和结束位置。

final List<PartitionInfo> partitionInfos = kafkaConsumer.partitionsFor("theTopicName");
final List<TopicPartition> topicPartitions = new ArrayList<>();
partitionInfos.forEach(partitionInfo -> topicPartitions.add(new TopicPartition("theTopicName", partitionInfo.partition())));
final List<OffsetRange> offsetRanges = new ArrayList<>();
kafkaConsumer.assign(topicPartitions);
topicPartitions.foreach(topicPartition -> {
    long fromOffset = kafkaConsumer.position(topicPartition);
    kafkaConsumer.seekToEnd(Collections.singleton(topicPartition));
    long untilOffset = kafkaConsumer.position(topicPartition);
    offsetRanges.add(new OffsetRange(topicPartition.topic(), topicPartition.partition(), fromOffset, untilOffset));
});
return offsetRanges.toArray(new OffsetRange[offsetRanges.size()]);

【问题讨论】:

    标签: java apache-spark spark-streaming-kafka sparkcore


    【解决方案1】:

    如果您想控制流量,您可以考虑使用轮询而不是流式 API。这样一来,您就可以在达到目标后明确停止投票。

    也看看这个...

    https://github.com/dibbhatt/kafka-spark-consumer

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-06-22
      • 2017-02-06
      • 2018-09-22
      • 1970-01-01
      • 2021-05-22
      • 2020-09-03
      • 2019-07-30
      • 2019-06-11
      相关资源
      最近更新 更多