scala – Akka.Kafka – 警告消息 – 恢复分区
发布时间:2020-12-16 18:31:52 所属栏目:安全 来源:网络整理
导读:我正在继续调试消息,因为它恢复了所有主题的分区.如下.此消息在我的服务器上连续打印每毫秒. 08:44:34.850 [default-akka.kafka.default-dispatcher-10] DEBUG o.a.k.clients.consumer.KafkaConsumer - Resuming partition test222-708:44:34.850 [default-a
我正在继续调试消息,因为它恢复了所有主题的分区.如下.此消息在我的服务器上连续打印每毫秒.
08:44:34.850 [default-akka.kafka.default-dispatcher-10] DEBUG o.a.k.clients.consumer.KafkaConsumer - Resuming partition test222-7 08:44:34.850 [default-akka.kafka.default-dispatcher-10] DEBUG o.a.k.clients.consumer.KafkaConsumer - Resuming partition test222-6 08:44:34.850 [default-akka.kafka.default-dispatcher-10] DEBUG o.a.k.clients.consumer.KafkaConsumer - Resuming partition test222-9 08:44:34.850 [default-akka.kafka.default-dispatcher-10] DEBUG o.a.k.clients.consumer.KafkaConsumer - Resuming partition test222-8 这个 val zookeeperHost = "localhost" val zookeeperPort = "9092" // Kafka queue settings val consumerSettings = ConsumerSettings(system,new ByteArrayDeserializer,new StringDeserializer) .withBootstrapServers(zookeeperHost + ":" + zookeeperPort) .withGroupId((groupName)) .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"latest") // Streaming the Messages from Kafka queue Consumer.committableSource(consumerSettings,Subscriptions.topics(topicName)) .map(msg => { consumed(msg.record.value) }) .runWith(Sink.ignore) 请帮助正确执行分区以停止DEBUG消息. 解决方法
似乎
reactive-kafka code在开始获取之前恢复每个分区:
consumer.assignment().asScala.foreach { tp => if (partitionsToFetch.contains(tp)) consumer.resume(java.util.Collections.singleton(tp)) else consumer.pause(java.util.Collections.singleton(tp)) } def tryPoll{...} checkNoResult(tryPoll(0)) 如果以前没有暂停分区,KafkaConsumer.resume方法是无操作的. (编辑:李大同) 【声明】本站内容均来自网络,其相关言论仅代表作者个人观点,不代表本站立场。若无意侵犯到您的权利,请及时与联系站长删除相关内容! |