【问题标题】:KafkaConsumer 0.10 Java API error message: No current assignment for partitionKafkaConsumer 0.10 Java API 错误消息:分区没有当前分配
【发布时间】:2017-04-21 20:58:39
【问题描述】:

我正在使用 KafkaConsumer 0.10 Java api。我想从特定的分区和特定的偏移量中消费。我查了一下,发现有一个 seek 方法,但是它抛出了一个异常。有人有类似的用例或解决方案吗?

代码:

KafkaConsumer<String, byte[]> consumer = new KafkaConsumer<>(consumerProps);
consumer.seek(new TopicPartition("mytopic", 1), 4);

例外

java.lang.IllegalStateException: No current assignment for partition mytopic-1
    at org.apache.kafka.clients.consumer.internals.SubscriptionState.assignedState(SubscriptionState.java:251)
    at org.apache.kafka.clients.consumer.internals.SubscriptionState.seek(SubscriptionState.java:276)
    at org.apache.kafka.clients.consumer.KafkaConsumer.seek(KafkaConsumer.java:1135)
    at xx.xxx.xxx.Test.main(Test.java:182)

【问题讨论】:

    标签: java kafka-consumer-api


    【解决方案1】:

    在您可以seek() 之前,您首先需要将subscribe() 到一个主题 assign() 将一个主题分区到消费者。还要记住,subscribe()assign() 是惰性的——因此,您还需要对 poll() 进行“虚拟调用”,然后才能使用 seek()

    注意:从 Kafka 2.0 开始,新的 poll(Duration timeout) 是异步的,并且不能保证当 poll 返回时您有完整的分配。因此,您可能需要在使用 seek() 和再次使用 poll 之前检查您的分配以刷新分配。 (详见KIP-266

    如果使用subscribe(),则使用组管理:因此,您可以使用相同的group.id启动多个消费者,主题的所有分区将自动平均分配给组内的所有消费者(每个分区将获得分配给组中的单个消费者)。

    如果要读取特定分区,需要通过assign()手动分配。这让你可以做任何你想做的任务。

    顺便说一句:KafkaConsumer 有一个很长的详细类 JavaDoc,包括示例。值得一读。

    【讨论】:

    • 谢谢。它工作:) 与 assign() 和 seek() 的组合
    • 我认为你的意思是 group.id 而不是 application.id
    • 在这里回答太多#KafkaStream 问题...感谢您指出@automaticgiant
    • @MatthiasJ.Sax "另外请记住,subscribe() 和 assign() 是惰性的——因此,您还需要对 poll() 进行“虚拟调用”,然后才能使用寻找()。”订阅和分配的惰性类型如何影响代码?如果我不使用像poll(0) 这样的虚拟投票怎么办?
    • 我猜 assign() 没问题。但是对于subscribe(),只要你不打电话给poll(),你就没有加入消费者组,也没有分配任何分区。因此,您最终会遇到问题中所示的异常。
    【解决方案2】:

    如果您不想使用 poll() 并检索地图记录,请更改偏移量本身。 卡夫卡 0.11 版 试试这个:

    ...
    props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");    
    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);    
    consumer.subscribe(Arrays.asList("Test_topic1", "Test_topic2"));
    List<TopicPartition> partitions =consumer.partitionsFor("Test_topic1").stream().map(part->{TopicPartition tp = new TopicPartition(part.topic(),part.partition()); return tp;}).collect(Collectors.toList());
    Field coordinatorField = consumer.getClass().getDeclaredField("coordinator"); 
    coordinatorField.setAccessible(true);    
    
    ConsumerCoordinator coordinator = (ConsumerCoordinator)coordinatorField.get(consumer);
    coordinator.poll(new Date().getTime(), 1000);//Watch out for your local date and time settings
    consumer.seekToBeginning(partitions); //or other seek
    

    对协调员事件进行投票。这确保了协调者是已知的并且消费者已经加入了组(如果它正在使用组管理)。如果启用,这也会处理定期偏移提交。

    【讨论】:

      【解决方案3】:

      请使用consumer。assign与consumer。seek而不是consumer。subscribe

      经过这些修改后,就可以正常执行了。

      【讨论】:

        猜你喜欢
        • 2019-08-14
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-06-28
        • 1970-01-01
        • 1970-01-01
        • 2016-10-15
        • 2019-04-15
        相关资源
        最近更新 更多