【问题标题】:Storm 0.10.0 reuse a topology design?Storm 0.10.0 重用拓扑设计?
【发布时间】:2015-12-23 04:02:28
【问题描述】:

下面的设计可以在Storm中完成吗?

让我们以下面的 wordcount 为例 https://github.com/nathanmarz/storm-starter/blob/master/src/jvm/storm/starter/WordCountTopology.java 我正在将单词生成器 spout 更改为文件阅读器 spout

这个字数统计拓扑的设计是 1.Spout读取文件并逐行创建句子 2.螺栓将句子拆分为单词 3. 用螺栓添加唯一词并给出一个词及其对应的计数

因此,拓扑在某种程度上描述了文件需要采用的流程来计算它所具有的唯一词。

如果我有两个文件,文件 1 和文件 2,一个应该能够调用相同的拓扑并创建此拓扑的两个实例以运行相同的字数。

为了跟踪字数统计是否确实已完成,一旦文件处理完毕,字数统计拓扑的实例应处于已完成状态。

在 Storm 的当前设计中,我发现 Topology 是实际的实例,所以它就像一个任务。

一个人需要使用不同的拓扑名称进行两次不同的调用,例如

对于文件 1 StormSubmitter.submitTopology("WordCountTopology1", conf,builder.createTopology());

对于文件 2 StormSubmitter.submitTopology("WordCountTopology2", conf,builder.createTopology());

更不用说使用storm客户端同样上传jar

storm jar stormwordcount-1.0.0-jar-with-dependencies.jar com.company.WordCount1Main.App "server" "filepath1"

storm jar stormwordcount-1.0.0-jar-with-dependencies.jar com.company.WordCount2Main.App "server" "filepath2"

另一个问题是处理文件后拓扑未完成。在我们对拓扑发出 kill 之前,它们一直都处于活动状态

storm kill "WordCountTopology"

我了解,在消息来自消息队列(如 Kafka)的流媒体世界中,消息没有尽头,但在实体/消息已修复的文件世界中,这与它有什么关系。

是否有执行以下操作的 API?

//创建拓扑,这是一次使用storm上传各自的jar StormSubmitter.submitTopology("WordCountTopology", conf,builder.createTopology());

上传后,应用程序代码只需使用 agruments 实例化拓扑 //创建拓扑实例并提供状态跟踪器 JobTracker tracker = StormSubmitter.runTopology("WordCountTopology", conf, args);

//可以在Storm中查询当前作业是否完成 JobStatus status = StormSubmitter.getTopologyStatus(conf, tracker);

【问题讨论】:

    标签: apache-storm


    【解决方案1】:

    重复使用相同的拓扑有两种可能:

    1) 为您的文件 spout 使用构造函数参数,并使用不同的参数两次实例化相同的拓扑:

    private StormTopology createMyTopology(String filename) {
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout("File Spout", new FileSpout(filename));
        // add further spouts and bolts etc.
        return builder.createTopology();
    }
    
    public static void main(String[] args) {
        String file1 = "/path/to/file1";
        String file2 = "/path/to/file2";
        Config c = new Config();
        if(useFile1) {
            StormSubmitter.submitTopology("T1", c, createMyTopology(file1));
        } else {
            StormSubmitter.submitTopology("T1", c, createMyTopology(file2));
        }
    }
    

    2) 作为替代方案,您可以在 open() 方法中配置文件 spout。

    public class FileSpout extends IRichSpout {
        @Override
        public void open(Map conf, ...) {
            String filenmae = (String)conf.get("FILENAME");
            // ...
        }
        // other methods omitted
    }
    
        public static void main(String[] args) {
        String file1 = "/path/to/file1";
        String file2 = "/path/to/file2";
        Config c = new Config();
        if(useFile1) {
            c.put("FILENAME", file1);
        } else {
            c.put("FILENAME", file2);
        }
    
        // assembly topology...
    
        StormSubmitter.submitTopology("T", c, builder.createTopology());
    }
    

    第二个问题:Storm 中没有 API 可以自动终止拓扑。您可以使用 TopologyInfo 并监控 spout 发出的元组的数量。如果它在一段时间内没有改变,你可以假设整个文件被读取然后终止拓扑。

    Config cfg = new Config();
    // set NIMBUS_HOST and NIMBUS_THRIFT_PORT in cfg
    Client client = NimbusClient.getConfiguredClient(cfg).getClient();
    TopologyInfo info = client.getTopologyInfo("topologyName");
    // get emitted tuples...
    client.killTopology("topologyName");
    

    【讨论】:

    • 感谢 Matthias,感谢提供一种可能的方法,随着文件数量的增加或必须使用相同的拓扑同时运行许多文件,而不是文件 1 或文件 2,这种顺序处理的设计将无济于事。我猜测文件消息也需要转换为流消息,并且在转换为流之前需要独立完成跟踪,在处理流之后需要更新跟踪以检查可能的消息结束。
    • 您的问题要求的内容与您的评论不同;)但是,我不确定您所说的“文件消息需要转换为流按摩”或“跟踪需要在转换之前独立完成”是什么意思成溪流”。 “需要更新跟踪以检查可能的消息结束”是什么意思?请更新您的问题并描述您的设置以及您想要完成的任务。现在,我有点迷茫……
    • 让我试一试我想要实现的目标。一个文件可以有 N 条消息,假设文件 1 有 100 条消息,文件 2 有 200 条消息。同样,可以有 M 个文件,每个文件都有 N 条消息。所有文件都需要使用相同的拓扑进行处理,因为所有文件中的所有消息都具有相同的处理要求。现在,由于 Storm 的设计,我必须为文件 1 、文件 2 ...等创建具有不同名称的相同拓扑。另一个问题似乎是如何在处理每个相应的文件后终止拓扑,例如 100 条消息后的文件 1 和 200 条消息后的文件 2。
    • cont..因为在文件 1 消息或文件 2 消息结束后拓扑不会消失,所以我无法找到消息的结尾。关于您的评论“获取发出的消息”以检查文件是否已处理,我如何以及在哪里编写此代码以检查处理了哪个文件?什么是 ack 和 fail 方法?如果文件 1 的 100 条消息和文件 2 的 200 条消息全部接收这两种方法,为什么在 ack 或 fail 方法中无法杀死拓扑。这样我就可以跟踪在每个拓扑实例中处理的文件的所有消息。
    • 您还可以在一个 spout 中一个接一个地处理多个文件。只需将整个文件列表(作为 List 或 String[] 提供给您的 spout)。如果您的文件之间没有任何依赖关系并且可以按任何顺序处理它们,您还可以并行化(有或没有并行性)您的 spout(或在单个拓扑中使用多个 spout)。关于“获取发出的消息”:此代码将转到一个额外的驱动程序,该程序监视正在运行的拓扑的进度,如果没有任何进展,则将其终止(假设没有进展表示所有消息都已处理)跨度>
    【解决方案2】:

    帖子中提到的字数统计拓扑并不能充分体现 Storm 的威力。由于 Storm 是一个流处理器,它需要一个流;时期。根据定义,文件是静态文件。我对 Storm 开发人员表示同情,即如何将一个简单的 hello world 用于展示拓扑概念和非流技术(如文件)的采用。所以对于我当时正在学习Storm的新手来说,如何使用示例进行开发是很难理解的。该示例只是展示 Storm 概念如何工作的一种方式,而不是文件将如何到达或需要处理的真实应用程序。

    以下是对其中一种解决方案的看法。

    由于拓扑一直在运行,因此它们可以计算任意时间段内的字数,即在一个文件内或跨所有文件的任意时间段内。

    为了允许不同的文件进入,我们需要一个流式喷口。所以很自然地,你需要一个 Kafka Message Broker 或类似的东西来接收流中的文件。根据文件的大小和消息代理设置的限制,即具有 1 MB 文件限制的 Kafka,我们可以选择将文件本身作为有效负载或文件的引用发送,在这种情况下,您需要一个分布式文件用于存储文件的系统,即 Hadoop DFS 或 NAS。

    然后我们使用 Kafka Spout 而不是 FileSpout 读取这些文件。

    我们现在有以下问题 1. 跨文件的字数统计 2. 每个文件的字数 3. 字数的运行状态,直到它被处理 4. 我们何时知道文件是否已处理或完成

    1. 跨文件字数统计

    使用提供的示例,这是示例目标的用例,因此如果我们继续流式传输文件并在每个文件中读取行,拆分单词并发送到其他螺栓,螺栓将独立计算单词它来自哪个文件。

    File1 一只快速的棕色狐狸跳了起来...... File2 仙界狐狸...

    字段分组 快的 棕色的 狐狸 ... 一次 之上 狐狸(不需要,因为它在文件 1 中) ...

    1. 每个文件的字数 为了做到这一点,我们现在需要将要附加到 fileId 的单词字段分组。因此,现在该示例需要更改为它拆分的每个单词都包含一个 fileId。 所以 File1 一只敏捷的棕色狐狸跳了起来…… File2 曾几何时一只狐狸...

    因此按单词分组的字段将是(取消干扰词)

    File1_quick File1_brown File1_fox

    File2_once File2_upon File2_fox

    1. 字数统计的运行状态,直到它被处理 由于所有这些计数都在螺栓的内存中,并且我们不知道 EoF,因此除非有人进入螺栓或我们定期将计数发送到另一个我们可以查询它的数据存储,否则无法获取状态。这正是我们需要做的,我们需要定期将内存中的 bolt 计数持久化到 hbase、elastic、mongo db 等数据存储中

    2. 我们何时知道文件是否已处理或完成 也许这是流媒体世界中最难回答的问题,基本上流处理器不知道流已经完成,因为从它的角度来看,流是进入的文件,它需要将每个文件分成单词并计算相应的螺栓。所以他们不知道在它到达每个演员之前或之后发生了什么。 整个事情需要由应用程序开发人员完成。 一种方法是在读取每个文件时计算总字数并发送消息 文件 1:总字数:1000 文件 2:总字数:2000

    现在,当我们计算字数并找到每个文件 File1_* 的不同单词时,在我们说文件完成之前,单个单词的计数和总单词应该匹配。所有这些都是我们需要编写的自定义逻辑,然后才能说它完整。

    所以本质上,Storm 提供了以多种方式进行流处理的框架。应用程序开发人员的工作是使用其拥有的设计进行开发,并根据用例实现自己的逻辑。它没有提供开箱即用的应用程序用例或一个很好的参考实现,我认为我们需要将其构建为它不是商业产品,并且依赖于社区来支持。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-05-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-30
      相关资源
      最近更新 更多