【发布时间】:2016-12-25 17:45:40
【问题描述】:
在使用 KafkaSpout 和几个 Bolt 开发和执行我的 Storm (1.0.1) 拓扑后,我注意到即使在拓扑空闲时也会产生巨大的网络流量(Kafka 上没有消息,bolts 中没有处理)。所以我开始逐条注释掉我的拓扑结构,以便找到原因,现在我的主目录中只有 KafkaSpout:
....
final SpoutConfig spoutConfig = new SpoutConfig(
new ZkHosts(zkHosts, "/brokers"),
"files-topic", // topic
"/kafka", // ZK chroot
"consumer-group-name");
spoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
spoutConfig.startOffsetTime = OffsetRequest.LatestTime();
topologyBuilder.setSpout(
"kafka-spout-id,
new KafkaSpout(config),
1);
....
当这个(无用的)拓扑执行时,即使是在本地模式下,即使是第一次,网络流量总是会增长很多:我看到(在我的活动监视器中)
- 平均每秒接收 432 KB 的数据
- 几个小时后,拓扑开始运行(空闲),接收到的数据为 1.26GB,发送的数据为 1GB
(重要提示:Kafka 没有在集群中运行,单个实例运行在具有单个主题和单个分区的同一台机器上。我刚刚在我的机器上下载了 Kafka,启动它并创建了一个简单的主题。当我把主题中的一条消息,拓扑中的所有内容都可以正常工作)
很明显,原因在KafkaSpout.nextTuple()方法(下),但我不明白为什么,在Kafka中没有任何消息,我应该有这样的流量。有什么我没有考虑到的吗?这是预期的行为吗?我看了看 Kafka 日志,ZK 日志,什么都没有,我已经清理了 Kafka 和 ZK 数据,什么都没有,仍然是相同的行为。
@Override
public void nextTuple() {
List<PartitionManager> managers = _coordinator.getMyManagedPartitions();
for (int i = 0; i < managers.size(); i++) {
try {
// in case the number of managers decreased
_currPartitionIndex = _currPartitionIndex % managers.size();
EmitState state = managers.get(_currPartitionIndex).next(_collector);
if (state != EmitState.EMITTED_MORE_LEFT) {
_currPartitionIndex = (_currPartitionIndex + 1) % managers.size();
}
if (state != EmitState.NO_EMITTED) {
break;
}
} catch (FailedFetchException e) {
LOG.warn("Fetch failed", e);
_coordinator.refresh();
}
}
long diffWithNow = System.currentTimeMillis() - _lastUpdateMs;
/*
As far as the System.currentTimeMillis() is dependent on System clock,
additional check on negative value of diffWithNow in case of external changes.
*/
if (diffWithNow > _spoutConfig.stateUpdateIntervalMs || diffWithNow < 0) {
commit();
}
}
【问题讨论】:
-
您是否也在storm-user邮件列表中提出过您的问题?
-
是的,我之前在这里尝试过,因为我不确定这件事更多的是在 Kafka 还是 Storm 方面