这是一个创建于 2340 天前的主题,其中的信息可能已经有所发展或是发生改变。
消费 kafka 数据时刚开始可以稳定跑一会,但是过不了几分钟就跑出此异常程序中断,java.lang.IllegalStateException: No current assignment for partition
我认为可能有问题的代码是 subscribe(),看网上有说用 Assign(),但是那样需要指定 partition,下面是我现在的代码:
val lineDStream: InputDStream[ConsumerRecord[Object, Object]] = KafkaUtils.createDirectStream(
ssc,
LocationStrategies.PreferConsistent,
ConsumerStrategies.Subscribe(kafkaTopics, kafkaParams)
)
如果有哪位知道解决方法,请指教,非常感谢!
第 1 条附言 · 2018-07-31 13:00:26 +08:00
问题已解决,是因为我在集群上跑着消费程序,本地也在用相同的消费代码测试,结果就出现了同一个 groupID 在同一时刻多次消费同一个 topic,引发 offset 记录问题。