【问题标题】:kafka spout is not emitting datakafka spout 没有发出数据
【发布时间】:2015-08-06 21:04:16
【问题描述】:

我正在尝试将 Kafka 与 Storm 集成。我正在使用 Kafka Spout 从 Kafka 主题中检索数据并将其提供给 Storm Bolt 以进行进一步处理。我能够成功提交拓扑,但 Spout 没有发出任何信息data.It 也不会引发任何错误。我对 Kafka 和 Storm 很陌生。所以,我无法找到这个问题背后的原因。请提出修改建议。提前致谢!

我的拓扑:

public class TopologyMain {

 private static final String SENTENCE_SPOUT_ID = "kafka-sentence-spout";


public static void main(String[] args) throws InterruptedException, AlreadyAliveException, InvalidTopologyException {
    int numSpoutExecutors = 1;


    TopologyBuilder builder = new TopologyBuilder();

    builder.setSpout(SENTENCE_SPOUT_ID, buildKafkaSentenceSpout(), numSpoutExecutors);
    builder.setBolt("word-normalizer", new WordNormalizer())
        .shuffleGrouping(SENTENCE_SPOUT_ID);
    builder.setBolt("word-counter", new WordCounter(),2)
        .shuffleGrouping("word-normalizer");

    //Configuration
    Config conf = new Config();
    conf.setDebug(false);
    //Topology run
    conf.put(Config.TOPOLOGY_MAX_SPOUT_PENDING, 1);
    conf.put(Config.NIMBUS_HOST, "192.168.1.229");
    conf.put(Config.NIMBUS_THRIFT_PORT, 6627);
    System.setProperty("storm.jar", "/home/ubuntu/st/stIn/target/storm-wc.jar");
    StormSubmitter.submitTopology("Count-Word-Topology", conf,builder.createTopology());

}



 private static KafkaSpout buildKafkaSentenceSpout() {
      BrokerHosts hosts = new ZkHosts("localhost:2181");
      SpoutConfig spoutConfig = new SpoutConfig(hosts, "test", "/acking-kafka-sentence-spout", "acking-sentence-spout");
      spoutConfig.forceFromStart = true;
      spoutConfig.startOffsetTime = kafka.api.OffsetRequest.EarliestTime();
      return new KafkaSpout(spoutConfig);
    }
 }

【问题讨论】:

  • 您可以使用控制台脚本来消费主题吗?这个拓扑在本地工作吗?您可以从集群中的storm-starter 运行任何示例拓扑吗?
  • 是的,我可以使用控制台脚本来处理主题。拓扑在本地运行良好。此外,我可以在 spout 是文本文件的集群上运行不同的风暴拓扑。
  • 如果你的 zookeeper 运行在与Config.NIMBUS_HOST 相同的 ip 中,请尝试将 localhost:2181 更改为 :2181 看看是否有帮助

标签: java apache-kafka apache-storm kafka-consumer-api


【解决方案1】:

我将项目的 Maven 依赖项中的所有 jar 明确复制到了storm库,一切正常。我还将storm jar(用于提交拓扑的jar)复制到storm/lib。

【讨论】:

  • 我是新手,面临同样的问题,你的意思是我的项目的 Maven 依赖项中的 jars,我提交的 jars 吗?
  • @user3188912 将项目的 maven 依赖项中存在的所有 jar 复制到storm lib..也复制storm topology jar 并设置此属性。 System.setProperty("storm.jar", "location of the topology jar");
  • 检查您的 kafka spout 在哪个端口上运行..然后检查相应的工作日志..这将为您提供 spout 失败的确切原因。
  • 使用新的稳定二进制文件总是更好..我建议你使用storm 0.9.4以后..因为其中修复了许多常见错误..如果你是没关系是否使用kafka ..您需要在storm lib中存在所有jar以及topology jar来运行代码..如果您不想这样做..然后创建一个阴影jar或fat jar并提交使用它的拓扑结构。
  • 我使用的是 Ubuntu 14.0.4。我认为您的问题与 Ubuntu 版本无关。将您的问题发布在storm 标签下。我一定可以帮助您,其他人也可以接受从我们的讨论中受益。
猜你喜欢
  • 1970-01-01
  • 2017-02-23
  • 2016-11-03
  • 2014-12-03
  • 2016-11-12
  • 2018-08-14
  • 1970-01-01
  • 1970-01-01
  • 2013-06-24
相关资源
最近更新 更多