【发布时间】: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),而不是在其运行期间。谢谢您的帮助。
【问题讨论】:
-
对于任何正在寻找“如何维护 ZooKeeper 的偏移值”的人,下面的链接用 python 代码解释。 stackoverflow.com/questions/44110027/…
标签: apache-spark apache-kafka spark-streaming offset spark-streaming-kafka