【问题标题】:KafkaSpout (idle) generates a huge network trafficKafkaSpout(空闲)产生巨大的网络流量
【发布时间】: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 方面

标签: apache-kafka apache-storm


【解决方案1】:

在 nextTuple() 方法中休眠一秒(1000ms),现在观察流量,例如,

@Override
public void nextTuple() {
   try {
       Thread.sleep(1000);
   } catch(Exception ex){
        log.error("Ëxception while sleeping...",e);
   }
   List<PartitionManager> managers = _coordinator.getMyManagedPartitions();
   for (int i = 0; i < managers.size(); i++) {
     ...
     ...
     ...
     ...
}

原因是,kafka consumer 是在pull methodology 的基础上工作的,也就是说,consumer 会从 kafka broker 中拉取数据。因此,从消费者的角度来看(Kafka Spout)将不断地向 kafka 代理发出一个获取请求,即TCP network request。因此,您将面临发送/接收数据包的大量统计数据。 Though the consumer doesn't consumes any message, pull request and empty response also will get account into network data packet sent/received statistics. 如果您的睡眠时间较长,您的网络流量将会减少。还有一些network related configurations 用于经纪人和消费者。对配置进行研究可能会对您有所帮助。希望对你有帮助。

【讨论】:

  • 我以与您相同的结论结束:我创建了一个 KafkaSpout 装饰器,它在调用 decoratee.nextTuple() 之前具有可配置的睡眠时间。设置 10 毫秒的延迟可将网络流量减少 10 倍(!)。谢谢回答
【解决方案2】:

您的螺栓是否收到消息?您的螺栓是否继承了 BaseRichBolt ?

在 Kafaspout 中注释掉 m.fail(id.offset) 行并检查一下。如果你的 bolt 没有确认,那么你的 spout 会假设消息失败并尝试重播相同的消息。

public void fail(Object msgId) {
        KafkaMessageId id = (KafkaMessageId) msgId;
        PartitionManager m = _coordinator.getManager(id.partition);
        if (m != null) {
            //m.fail(id.offset);
        }

还可以尝试将 nextTuple() 暂停几毫秒并检查一下。

如果有帮助请告诉我

【讨论】:

  • 您好,非常感谢您的回复。不幸的是,我注释掉了我的螺栓……在拓扑中只有 KafkaSpout,我没有向 Kafka 发送任何消息。我正在调试 KafkaSpout,它只是(正确地)从 Kafka 获取零消息负载......但似乎这件事(在 nextTuple 中连续完成)会导致如此巨大的流量(我不明白为什么)跨度>
  • 你可以玩这个东西并检查一下 spoutConfig.startOffsetTime = OffsetRequest.LatestTime(); // -1, -2 kafkaConfig.forceStartOffsetTime(-1)...(-2) 用于最新和最早的偏移量..
  • 不确定我真的得到了你。但是在 ZK 中没有记录消费者状态,因为 Kafka Spout 从未消费过消息,所以我看到(瞬态)偏移设置为 -1。如果我取消注释螺栓,然后在主题中放入一条消息,则整个元组树都会被处理,之后我可以看到 ZK 中记录的消费者状态(偏移量 = 1)。所以从功能的角度来看,一切都按预期工作......只有那个该死的交通拥堵仍然存在
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-03-26
  • 2023-03-24
  • 1970-01-01
  • 2023-02-03
  • 1970-01-01
  • 2014-07-10
相关资源
最近更新 更多