【问题标题】:How to run WordCountTopology from storm-starter in Intellij如何在 Intellij 中从storm-starting 运行WordCount Topology
【发布时间】:2015-08-13 10:36:55
【问题描述】:

我已经与 Storm 合作了一段时间,但想开始开发。正如建议的那样,我使用的是 IntelliJ(到目前为止,我使用的是 Eclipse,并且只针对 Java API 编写拓扑)。

我也在看 https://github.com/apache/storm/tree/master/examples/storm-starter#intellij-idea

此文档不完整。我无法首先在 Intellij 中运行任何东西。我可以弄清楚,我需要删除storm-core依赖的范围(在storm-starter pom.xml中)。 (在这里找到:storm-starter with intellij idea,maven project could not find class

之后我就可以构建项目了。我也可以在 IntelliJ 中毫无问题地运行 ExclamationTopology。但是,WordCountTopology 失败。

首先我收到以下错误:

java.lang.RuntimeException: backtype.storm.multilang.NoOutputException: 到子进程的管道似乎坏了!没有输出读取。 序列化程序异常: 回溯(最近一次通话最后): 文件“splitsentence.py”,第 16 行,在 进口风暴 ImportError: 没有名为storm的模块

更新:无需安装python-storm 即可使其工作

我能够通过以下方式解决它:apt-get install python-storm(来自 StackOverflow)

但是,我不会说 Python,我想知道问题出在哪里以及为什么我可以这样解决它。只是想更深入地了解它。也许有人可以解释一下。

不幸的是,我现在遇到了另一个错误:

java.lang.RuntimeException: backtype.storm.multilang.NoOutputException: 到子进程的管道似乎坏了!没有输出读取。 序列化程序异常: 回溯(最近一次通话最后): 文件“splitsentence.py”,第 18 行,在 类SplitSentenceBolt(storm.BasicBolt): AttributeError: 'module' 对象没有属性 'BasicBolt'

我在 Internet 上没有找到任何解决方案。在dev@storm.apache.org 询问也无济于事。我提出以下建议:

我认为总是假设拓扑总是会通过storm-命令行调用。因此工作目录将是 ${STORM-INSTALLATION}/bin/storm 由于storm.py 在这个目录中,splitSentence.py 将能够找到storm 模块。您可以将工作目录设置为storm.py所在的路径,然后尝试。如果可行,我们可以稍后将其添加到文档中

但是,改变工作目录并没有解决问题。

由于我不熟悉 Python,而且我是 IntelliJ 的新手,所以我现在被困住了。因为ExclamationTopology 运行,我猜我的基本设置是正确的。

我做错了什么?完全可以在 IntelliJ 中运行 WordcountTopology in LocalCluster 吗?

【问题讨论】:

    标签: python intellij-idea apache-storm


    【解决方案1】:

    不幸的是,AFAIK 如果没有打包文件,您将无法使用 LocalCluster 运行多语言功能。

    ShellProcess 依赖于 TopologyContext 的 codeDir,由 supervisor 使用。 Worker 被序列化为stormcode.ser,但多语言文件应该被提取到序列化文件之外,以便python/ruby/node/etc 可以加载它。

    使用分发模式很容易做到这一点,因为总是有用户提交的jar,主管可以知道这是用户提交的。

    但是用本地模式做到这一点并不容易,因为主管无法知道用户提交的jar,用户可以在不打包的情况下将拓扑运行到本地模式。

    因此,本地模式下的主管从类路径中的每个 jar(以“jar”结尾)中查找资源目录(“resources”),并将第一次出现的地方复制到 codeDir。

    storm jar 将用户拓扑 jar 放置到类路径的第一个,因此它可以毫无问题地运行。

    所以通常情况下,ShellProcess 找不到“splitsentence.py”是很自然的。也许你的工作目录或 PYTHONPATH 成功了。

    【讨论】:

    • 设置 PYTHONPATH 可以解决问题! SplitSentence pythonSplit = new SplitSentence(); Map env = new HashMap(); env.put("PYTHONPATH", "/home/mjsax/workspace_storm/storm/storm-multilang/python/src/main/resources/resources/"); pythonSplit.setEnv(env); builder.setBolt("split",pythonSplit, 8).shuffleGrouping("spout");
    • 是的,如果找不到storm-starter的多语言实现,你也可以将它们添加到PYTHONPATH中。
    • 风暴启动器中不包含它。并不是storm.py 的每个版本都有效。该代码包含 3 个不同版本的 5 个文件。无法弄清楚为什么有不同的版本。你知道吗?
    • bin/storm.py 是 Storm 组件的启动器/shell 脚本。 (storm-dist/binary/target/apache-storm-0.11.0-SNAPSHOT/bin/storm.py是拷贝版)
    • 谢谢。这是常规的 Maven 构建过程。我很清楚这就是同一文件的多个副本的原因。向后兼容性——我明白了。有道理。也许它应该记录在某个地方...... ;)
    【解决方案2】:

    我遇到了类似的问题,不是示例拓扑,而是我自己使用 Python 螺栓。

    还遇到了“AttributeError: 'module' object has no attribute 'BasicBolt'”异常——在本地模式下和提交到集群时。

    这方面的资源很少,我发现了你的问题,很少有人讨论这个问题。

    如果其他人有同样的问题: 确保在 pom 文件中包含正确的 Maven“multilang-python”依赖项。这会将正确的运行时依赖项打包到运行拓扑所需的 JAR 文件中。

    【讨论】:

    • 对不起,我遇到了同样的问题,但无法解决,我怎么知道我应该使用正确的 maven 吗?
    • 如果我没记错的话,您必须将 multilang-python 依赖项上的版本号与主风暴依赖项匹配
    【解决方案3】:

    我设法在我的 virtualbox 上运行它,风暴版本 1.2.2:

    只需下载https://github.com/apache/storm/blob/master/storm-multilang/python/src/main/resources/resources/storm.py并将其放入您想要的任何文件夹中,例如: /apache-storm-1.2.2/examples/storm-starter/multilang/resources/ ,然后更改主要功能:

    public static void main(String[] args) throws Exception {
    
        SplitSentence pythonSplit = new SplitSentence();
        Map env = new HashMap();
        env.put("PYTHONPATH", "/apache-storm-1.2.2/examples/storm-starter/multilang/resources/");
        pythonSplit.setEnv(env);
    
        TopologyBuilder builder = new TopologyBuilder();
    
        builder.setSpout("spout", new RandomSentenceSpout(), 5);
    
        builder.setBolt("split",pythonSplit, 8).shuffleGrouping("spout");
        builder.setBolt("count", new WordCount(), 12).fieldsGrouping("split", new Fields("word"));
    
        Config conf = new Config();
        conf.setDebug(true);
    
        if (args != null && args.length > 0) {
          conf.setNumWorkers(3);
    
          StormSubmitter.submitTopologyWithProgressBar(args[0], conf, builder.createTopology());
        }
        else {
          conf.setMaxTaskParallelism(3);
    
          LocalCluster cluster = new LocalCluster();
          cluster.submitTopology("word-count", conf, builder.createTopology());
    
          Thread.sleep(600000);
    
          cluster.shutdown();
        }
      }
    

    完整的说明可以在我的博客上找到,其中包括在本地模式和本地集群模式下运行时遇到的其他问题:https://lyhistory.com/storm/

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-02-07
      • 2016-09-14
      • 2023-03-29
      • 1970-01-01
      • 2016-03-23
      • 1970-01-01
      • 2018-01-29
      • 1970-01-01
      相关资源
      最近更新 更多