【发布时间】:2018-11-30 13:38:36
【问题描述】:
我正在测试简单的拓扑来检查 kafka spout 的性能。 它包含确认每个元组的 kafka spout 和 Bolt。 螺栓执行方法:
public void execute(Tuple input) {
collector.ack(input);
}
拓扑如下所示:
protected void configureTopology(TopologyBuilder topologyBuilder) {
configureKafkaCDRSpout(topologyBuilder);
configureKafkaSpoutBandwidthTesterBolt(topologyBuilder);
}
private void configureKafkaCDRSpout(TopologyBuilder builder) {
KafkaSpout kafkaSpout = new KafkaSpout(createKafkaCDRSpoutConfig());
int spoutCount = Integer.valueOf(topologyConfig.getProperty("kafka.cboss.cdr.spout.thread.count"));
builder.setSpout(KAFKA_CDR_SPOUT_ID, kafkaSpout, spoutCount)
.setNumTasks(Integer.valueOf(topologyConfig.getProperty(KAFKA_CDR_SPOUT_NUM_TASKS)));
}
private SpoutConfig createKafkaCDRSpoutConfig() {
BrokerHosts hosts = new ZkHosts(topologyConfig.getProperty("kafka.zookeeper.broker.host"));
String topic = topologyConfig.getProperty("kafka.cboss.cdr.topic");
String zkRoot = topologyConfig.getProperty("kafka.cboss.cdr.zkRoot");
String consumerGroupId = topologyConfig.getProperty("kafka.cboss.cdr.consumerId");
SpoutConfig kafkaSpoutConfig = new SpoutConfig(hosts, topic, zkRoot, consumerGroupId);
kafkaSpoutConfig.scheme = new SchemeAsMultiScheme(new CbossCdrScheme());
kafkaSpoutConfig.ignoreZkOffsets = true;
kafkaSpoutConfig.fetchSizeBytes = Integer.valueOf(topologyConfig.getProperty("kafka.fetchSizeBytes"));
kafkaSpoutConfig.bufferSizeBytes = Integer.valueOf(topologyConfig.getProperty("kafka.bufferSizeBytes"));
return kafkaSpoutConfig;
}
public void configureKafkaSpoutBandwidthTesterBolt(TopologyBuilder topologyBuilder) {
SimpleAckerBolt b = new SimpleAckerBolt();
topologyBuilder.setBolt(SPOUT_BANDWIDTH_TESTER_BOLT_ID, b, Integer.valueOf(topologyConfig.getProperty(CFG_SIMPLE_ACKER_BOLT_PARALLELISM)))
.setNumTasks(Integer.valueOf(topologyConfig.getProperty(SPOUT_BANDWIDTH_TESTER_BOLT_NUM_TASKS)))
.localOrShuffleGrouping(KAFKA_CDR_SPOUT_ID);
}
其他拓扑设置:
topology.max.spout.pending=250
topology.executor.receive.buffer.size=1024
topology.executor.send.buffer.size=1024
topology.receiver.buffer.size=8
topology.transfer.buffer.size=1024
topology.acker.executors=1
我正在使用 1 个工人、1 个 Kafka Spout 和 1 个 Simple Acker Bolt 启动我的拓扑。 这就是我在风暴 UI 中得到的:
好吧,我在 10 分钟内得到了 1.5kk 个元组。螺栓容量约为 0,5。所以我的逻辑很简单:如果我双喷嘴和螺栓并行提示 - 我将获得双倍性能。 下一个测试是使用 1 个 worker 2 个 Kafka Spout、2 个 Simple Acker Bolt 和 topology.acker.executors=2。结果如下:
因此,随着并行提示的增加,我的性能会变得更差。为什么会发生?如何增加每秒处理的元组?实际上,任何 spout 并行提示大于 2 的测试都显示出比 1 spout executor 更差的结果。
我已经检查过了:
1)这不是卡夫卡的错。主题在 2 个代理上有 20 个分区。 4 个工人规模的拓扑并获得 x4 性能。
2)这不是服务器故障。服务器有 40 个内核和 32Gb RAM。在运行拓扑时,它消耗大约 1/8 的 CPU 和几乎没有 RAM。
3) 更改 topology.max.spout.pending 参数没有帮助。
4) 增加 Bolt 或 Acker 并行性提示也无济于事。
【问题讨论】:
-
你只用一个工人运行了两个测试,如果你增加了另一个工人怎么办?所以用两个工人运行第二个测试。
-
感谢您的回复,摩根。你说得对。越来越多的工人给了我成比例的结果。有 2 个工人 2 个喷口,我每秒的元组数加倍。但是这个测试的想法是为了衡量 1 名工人的最佳表现。我能得到的最好的结果是每 10 分钟 1,5kk 个元组,或每秒 2500 个元组。我想在具有 40 个内核、32GB RAM 和 10Gb/s 网络的节点上我可以做得更好。
-
1 个工人,但 40 个核心并没有真正的意义。无论如何,每个工作人员都是一个线程,因此这意味着您的服务器有能力容纳 40 个工作人员。您现在在单个核心服务器上拥有完全相同的性能。每个线程 2500 元组/秒并不是很好,但也不算太差。
-
Storm 文档说 1 个 worker 就是 1 个 JVM。所以没有理由在 1 个节点上运行多个 worker。
-
@f1sherox,我不一定同意这种说法。运行多个 worker 的原因之一是为了提高容错性。如果您只有 1 个工作人员,如果该 1 个工作人员发生故障,您的整个拓扑都会失败。如果您有 6 个工作人员,并且 1 个工作人员失败,则 5 个工作人员仍在运行。此外,一个 worker 只能属于一个拓扑,因此只有 1 个 worker 意味着您的 Storm 集群一次只能支持 1 个拓扑。