【发布时间】: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;
}
我们是否可以通过增加并行任务的数量或/和避免在完成批处理时出现饥饿来提高性能?
【问题讨论】: