【问题标题】:Storm+Kafka not parallelizing as expectedStorm+Kafka 未按预期并行化
【发布时间】:2019-07-30 05:24:12
【问题描述】:

我们遇到了关于单个拓扑内的任务并行性的问题。我们无法获得良好、流畅的处理速度。

我们正在使用 Kafka 和 Storm 构建具有不同拓扑的系统,其中数据按照使用 Kafka 主题连接的拓扑链进行处理。

我们正在使用 Kafka 1.0.0 和 Storm 1.2.1。

负载很少,每天大约 2000 条消息,但是每个任务可能需要相当长的时间。特别是一种拓扑可能需要不同的时间来处理每个任务,通常在 1 到 20 分钟之间。如果按顺序处理,吞吐量不足以处理所有传入的消息。所有拓扑和 Kafka 系统都安装在一台机器上(16 核,16 GB RAM)。

由于消息是独立的并且可以并行处理,我们正在尝试使用 Storm 并发能力来提高吞吐量。

为此,拓扑已配置如下:

  • 4 名工人
  • 并行提示设置为 10
  • 从 Kafka 读取时的消息大小足以在每条消息中读取大约 8 个任务。
  • Kafka 主题使用复制因子 = 1 和分区 = 10。

通过此配置,我们在此拓扑中观察到以下行为。

  • Storm 拓扑从 Kafka 批量读取大约 7-8 个任务(任务大小不固定),最大消息大小为 128 kB。
  • 同时计算大约 4-5 个任务。工作在工人之间或多或少是平均分配的。一些工作人员负责 1 个任务,另一些工作人员负责 2 个任务并同时处理它们。
  • 随着任务的完成,剩余的任务开始。
  • 当只有 1-2 个任务需要处理时,就会出现饥饿问题。其他工作人员等待所有任务完成。
  • 当所有任务完成后,确认消息并发送到下一个拓扑。
  • 从 Kafka 读取一个新批次并重新开始该过程。

我们有两个主要问题。首先,即使有 4 个 worker 和 10 个并行提示,也只能启动 4-5 个任务。其次,当有待处理的工作时,不会再启动批次,即使它只是 1 个任务。

这不是没有足够的工作要做的问题,因为我们尝试在开始时插入 2000 个任务,所以有很多工作要做。

我们曾尝试增加参数“maxSpoutsPending”,期望拓扑会同时读取更多批次并将它们排队,但似乎它们正在内部流水线化,而不是同时处理。

使用以下代码创建拓扑:

private static StormTopology buildTopologyOD() {
    //This is the marker interface BrokerHosts.
    BrokerHosts hosts = new ZkHosts(configuration.getProperty(ZKHOSTS));
    TridentKafkaConfig tridentConfigCorrelation = new TridentKafkaConfig(hosts, configuration.getProperty(TOPIC_FROM_CORRELATOR_NAME));

    tridentConfigCorrelation.scheme = new RawMultiScheme();
    tridentConfigCorrelation.fetchSizeBytes = Integer.parseInt(configuration.getProperty(MAX_SIZE_BYTES_CORRELATED_STREAM));

    OpaqueTridentKafkaSpout spoutCorrelator = new OpaqueTridentKafkaSpout(tridentConfigCorrelation);

    TridentTopology topology = new TridentTopology();

    Stream existingObject = topology.newStream("kafka_spout_od1", spoutCorrelator)
            .shuffle()
            .each(new Fields("bytes"), new ProcessTask(), new Fields(RESULT_FIELD, OBJECT_FIELD))
            .parallelismHint(Integer.parseInt(configuration.getProperty(PARALLELISM_HINT)));

    //Create a state Factory to produce outputs to kafka topics.
    TridentKafkaStateFactory stateFactory = new TridentKafkaStateFactory()
            .withProducerProperties(kafkaProperties)
            .withKafkaTopicSelector(new ODTopicSelector())
            .withTridentTupleToKafkaMapper(new ODTupleToKafkaMapper());

    existingObject.partitionPersist(stateFactory, new Fields(RESULT_FIELD, OBJECT_FIELD), new TridentKafkaUpdater(), new Fields(OBJECT_FIELD));

    return topology.build();
}

和配置创建为:

private static Config createConfig(boolean local) {
    Config conf = new Config();
    conf.setMaxSpoutPending(1); // Also tried 2..6
    conf.setNumWorkers(4);

    return conf;
}

我们是否可以通过增加并行任务的数量或/和避免在完成批处理时出现饥饿来提高性能?

【问题讨论】:

    标签: apache-kafka apache-storm


    【解决方案1】:

    我在 Nathan Marz 的 Storm-users 上找到了一个 old post,关于为 Trident 设置并行性:

    我建议使用“名称”函数来命名流的部分内容 以便 UI 向您显示哪些螺栓对应于哪些部分。

    Trident 将操作打包到尽可能少的螺栓中。此外, 它从不重新分区您的流,除非您已完成操作 明确涉及重新分区(例如 shuffle、groupBy、 partitionBy、全局聚合等)。三叉戟的这个属性 确保您可以控制事物的排序/半排序 被处理。所以在这种情况下, groupBy 之前的所有内容都必须 具有相同的并行性,否则 Trident 将不得不重新分区 流。既然你没有说你想要流 重新分区,它不能这样做。你可以获得不同的并行度 通过引入重新分区来为喷口与每个人的追随者 操作,像这样:

    stream.parallelismHint(1).shuffle().each(…).each(…).parallelismHint(3).groupBy(…);

    我认为您可能希望为您的 spout 以及您的 .each 设置 parallelismHint。

    关于同时处理多个批次,你说得对,这就是maxSpoutPending 在 Trident 中的用途。尝试在 Storm UI 中检查您的最大 spout 挂起值是否实际被拾取。还可以尝试为MasterBatchCoordinator 启用调试日志记录。您应该能够从该日志中判断多个批次是否同时在运行。

    当你说多个批次不并发处理时,你是指ProcessTask吗?请记住,Trident 的属性之一是状态更新是按批次排序的。如果你有例如maxSpoutPending=3 并且批次 1、2 和 3 正在运行,Trident 不会发出更多批次进行处理,直到写入批次 1,此时它将再发出一个。所以慢批次可以阻止发射更多,即使 2 和 3 被完全处理,它们也必须等待 1 完成并被写入。

    如果您不需要 Trident 的批处理和排序行为,您可以尝试使用常规的 Storm。

    更多附注,但您可能需要考虑从 storm-kafka 迁移到 storm-kafka-client。这个问题不重要,但是不做就升级到Kafka 2.x,而且在得到一堆状态之前更容易迁移。

    【讨论】:

    • 感谢您的评论。正如您所说,我也尝试设置 parallelismHint 并获得相同的结果。正在获取 maxSpoutPending 的值,但显然没有使用,或者至少没有像我预期的那样使用。我们使用 Trident 是因为,AFAIK,正常拓扑可能无法处理某些任务或处理不止一次。我不介意三叉戟阻止发射,如果它们正在被处理并在最后发射,但似乎它们甚至没有被处理。
    • 在这种情况下,我不确定问题出在哪里。考虑试试storm-user邮件列表,那里可能有人有调整Trident拓扑的经验storm.apache.org/getting-help.html
    • 好的,我试试!谢谢
    猜你喜欢
    • 1970-01-01
    • 2021-04-08
    • 2022-01-21
    • 2019-09-23
    • 2012-11-14
    • 1970-01-01
    • 2020-05-08
    • 2017-07-19
    • 2020-12-06
    相关资源
    最近更新 更多