【问题标题】:Spark Structured Streaming Kafka Offset ManagementSpark 结构化流式处理 Kafka 偏移管理
【发布时间】:2019-10-03 05:44:44
【问题描述】:

我正在研究在 kafka 中存储用于 Spark 结构化流的 kafka 偏移量,就像它适用于 DStreams stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges) 一样,我正在寻找相同的东西,但适用于结构化流。 是否支持结构化流式传输?如果是,我该如何实现?

我知道使用 .option("checkpointLocation", checkpointLocation) 的 hdfs 检查点,但我对内置偏移管理完全感兴趣。

我希望 kafka 只在没有 spark hdfs 检查点的情况下在内部存储偏移量。

【问题讨论】:

    标签: apache-spark apache-kafka spark-structured-streaming spark-kafka-integration


    【解决方案1】:

    我正在使用在某处找到的这段代码。

    public class OffsetManager {
    
        private String storagePrefix;
    
        public OffsetManager(String storagePrefix) {
            this.storagePrefix = storagePrefix;
        }
    
        /**
         * Overwrite the offset for the topic in an external storage.
         *
         * @param topic     - Topic name.
         * @param partition - Partition of the topic.
         * @param offset    - offset to be stored.
         */
        void saveOffsetInExternalStore(String topic, int partition, long offset) {
    
            try {
    
                FileWriter writer = new FileWriter(storageName(topic, partition), false);
    
                BufferedWriter bufferedWriter = new BufferedWriter(writer);
                bufferedWriter.write(offset + "");
                bufferedWriter.flush();
                bufferedWriter.close();
    
            } catch (Exception e) {
                e.printStackTrace();
                throw new RuntimeException(e);
            }
        }
    
        /**
         * @return he last offset + 1 for the provided topic and partition.
         */
        long readOffsetFromExternalStore(String topic, int partition) {
    
            try {
    
                Stream<String> stream = Files.lines(Paths.get(storageName(topic, partition)));
    
                return Long.parseLong(stream.collect(Collectors.toList()).get(0)) + 1;
    
            } catch (Exception e) {
                e.printStackTrace();
            }
    
            return 0;
        }
    
        private String storageName(String topic, int partition) {
            return "Offsets\\" + storagePrefix + "-" + topic + "-" + partition;
        }
    
    }
    

    SaveOffset...在记录处理成功后调用,否则不存储偏移量。我使用 Kafka 主题作为源,所以我将起始偏移量指定为从 ReadOffsets 检索到的偏移量...

    【讨论】:

      【解决方案2】:

      “是否支持结构化流式传输?”

      不,结构化流不支持将偏移量提交回 Kafka,类似于使用 Spark Streaming (DStreams) 可以完成的操作。 Kafka specific configurations 上的 Spark Structured Streaming + Kafka 集成指南对此非常准确:

      “Kafka 源没有提交任何偏移量。”

      我已经在How to manually set groupId and commit Kafka offsets in Spark Structured Streaming 中写了一个更全面的答案。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-06-28
        • 1970-01-01
        • 2018-09-22
        • 2018-09-28
        • 2017-02-06
        • 2019-01-14
        • 2021-01-15
        • 2021-03-18
        相关资源
        最近更新 更多