【问题标题】:Spark Streaming kafka offset manageSpark Streaming kafka 偏移管理
【发布时间】:2018-09-22 17:19:23
【问题描述】:

我一直在做火花流工作,通过 kafka 消费和生产数据。我用的是directDstream,所以我必须自己管理offset,我们用redis来读写offset。现在有个问题,当我启动我的客户端时,我的客户端需要从redis中获取offset,而不是kafka中存在的offset本身。如何显示我编写我的代码?现在我已经在下面编写了我的代码:

   kafka_stream = KafkaUtils.createDirectStream(
    ssc,
    topics=[config.CONSUME_TOPIC, ],
    kafkaParams={"bootstrap.servers": config.CONSUME_BROKERS,
                 "auto.offset.reset": "largest"},
    fromOffsets=read_offset_range(config.OFFSET_KEY))

但我认为 fromOffsets 是 spark-streaming 客户端启动时的值(来自 redis),而不是在其运行期间。谢谢您的帮助。

【问题讨论】:

标签: apache-spark apache-kafka spark-streaming offset spark-streaming-kafka


【解决方案1】:

如果我理解正确,您需要手动设置偏移量。我就是这样做的:

from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils
from pyspark.streaming.kafka import TopicAndPartition

stream = StreamingContext(sc, 120) # 120 second window

kafkaParams = {"metadata.broker.list":"1:667,2:6667,3:6667"}
kafkaParams["auto.offset.reset"] = "smallest"
kafkaParams["enable.auto.commit"] = "false"

topic = "xyz"
topicPartion = TopicAndPartition(topic, 0)
fromOffset = {topicPartion: long(PUT NUMERIC OFFSET HERE)}

kafka_stream = KafkaUtils.createDirectStream(stream, [topic], kafkaParams, fromOffsets = fromOffset)

【讨论】:

  • 是的,我明白你的回答。但是,我认为参数“fromOffsets”意味着在客户端重新启动时获得偏移量。那么客户端在消费时如何获得偏移量?来自 kafka 本身或还是 fromOffsets?
  • 我不确定你的意思。当你开始一个流时,它必须从某个偏移量开始,然后从那里一直持续下去,直到它到达主题的结尾。参数“fromOffsets”的意思是:当你启动一个“createDirectStream”时,这意味着你不想从 ZooKeeper 读取偏移量,就像你使用“createStream”时一样,所以你需要自己提供。您是在问如何在不使用 ZooKeeper 的情况下获取偏移量?
  • 顺便说一句,我想问另一个问题,你如何调试你的python spark流,我使用“logger”输出有用的消息,但它在spark集群中不起作用(部署模式是客户端)
  • 调试有问题。通常我使用 jupyter 来预运行我的代码并在更友好的环境中对其进行测试。如果我必须在生产中进行调试,我会使用(非常“丑陋”)从 rdd 内部将日志写入 SQL 的做法
猜你喜欢
  • 2017-02-06
  • 2021-05-22
  • 2018-04-12
  • 2017-06-22
  • 1970-01-01
  • 2021-01-15
  • 2020-09-03
  • 1970-01-01
  • 2017-07-17
相关资源
最近更新 更多