【问题标题】:Kafka spout error "Consumer is not subscribed to any topics or assigned any partitions "Kafka spout 错误“消费者未订阅任何主题或分配任何分区”
【发布时间】:2017-10-28 02:23:40
【问题描述】:

我使用的是 Storm 1.1.0 版和 Kafka 0.10.1.2 版。

我正在按如下方式创建 Kafka-spout:

public KafkaSpout<String, String> getKafkaSpout() {
    String _kafkaBrokers = (String) props.get("bootstrap.servers");
    String _topic = (String) props.get("kafka.topic.name");
    String groupId = (String) props.get("group.id");
    int maxMsgSize = (int) props.get("fetch.message.max.bytes");
    String keySerializer = (String) props.get("key.serializer");
    String valueSerializer = (String) props.get("value.serializer");

    List<String>topics = new ArrayList<String>(`enter code here`);
    topics.add(_topic);

    return new KafkaSpout<String, String (KafkaSpoutConfig.builder(_kafkaBrokers, topics)
            .setFirstPollOffsetStrategy(FirstPollOffsetStrategy.UNCOMMITTED_EARLIEST)
            .setMaxUncommittedOffsets(100)
            .setProp(ConsumerConfig.GROUP_ID_CONFIG, groupId)
            .setProp(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG,maxMsgSize)
            .setProp("key.serializer",keySerializer)
            .setProp("value.serializer",valueSerializer)
            .build())
}

我收到下面提到的错误

java.lang.IllegalStateException: Consumer is not subscribed to any topics or assigned any partitions 
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:973) 
at org.apache.storm.kafka.spout.KafkaSpout.pollKafkaBroker(KafkaSpout.java:291) 
at org.apache.storm.kafka.spout.KafkaSpout.nextTuple(KafkaSpout.java:225) 
at org.apache.storm.daemon.executor$fn__9798$fn__9813$fn__9844.invoke(executor.clj:647) 
at org.apache.storm.util$async_loop$fn__555.invoke(util.clj:484) 
at clojure.lang.AFn.run(AFn.java:22) at java.lang.Thread.run(Thread.java:745) 

以及我在下面提到的其他依赖项在项目中的 maven 依赖项

<dependency>
    <groupId>org.apache.storm</groupId>
    <artifactId>storm-kafka-client</artifactId>
    <version>1.1.0.2.6.2.0-205</version>
</dependency>
<dependency>
    <groupId>org.apache.storm</groupId>
    <artifactId>storm-kafka</artifactId>
    <version>1.1.0.2.6.2.0-205</version>
</dependency>

【问题讨论】:

    标签: apache-kafka apache-storm


    【解决方案1】:

    我认为List&lt;String&gt;topics = new ArrayList&lt;String&gt;("enter code here"); 是您的问题?您可能需要在该列表中写下您的主题名称。

    你的依赖版本很奇怪,AFAIK Storm 没有发布任何带有这些版本字符串的东西。

    我也想知道为什么你需要storm-kafka-client,它适用于Kafka > 0.10集群,以及storm-kafka,它适用于较旧的Kafka集群(但我认为目前仍与最新的Kafka兼容)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-06-14
      • 2017-10-17
      • 2020-12-26
      • 2017-07-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-07-28
      相关资源
      最近更新 更多