You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Storm与RabbitMQ的单词计数拓扑无输出问题求助

Troubleshooting Your Storm Word Count Topology

Let's walk through the issues in your code step by step — there are several critical bugs preventing tuples from flowing through your topology and reaching the output queue:


1. RabbitMQSpout: Broken Collector Assignment

In your open() method, you're assigning the collector to itself instead of using the provided parameter:

// Wrong
_collector = _collector;
// Correct
_collector = spoutOutputCollector;

Without fixing this, your spout can never emit tuples because the collector reference is null.


2. RabbitMQManager: Invalid Message Consumption Logic

Your receive() method uses basicConsume(), which is asynchronous and returns a consumer tag string (not the actual message content). For Storm spouts, you need synchronous message fetching. Replace the receive() method with this implementation using basicGet():

public String receive(String queue) {
    try {
        if (!reopenConnectionIfNeeded()) {
            return null;
        }
        Channel channel = connection.createChannel();
        GetResponse response = channel.basicGet(queue, true);
        if (response != null) {
            String message = new String(response.getBody(), StandardCharsets.UTF_8);
            channel.close();
            return message;
        }
        channel.close();
    } catch (Exception e) {
        e.printStackTrace();
    }
    return null;
}

This will properly fetch individual messages from the RabbitMQ queue to emit from your spout.


3. SplitBolt: Class Name Mismatch & Undefined Variable

Your bolt class is named SplitBolt, but the constructor is SplitSentenceBolt(), and you're referencing an undefined SPACE variable. Fix the class name consistency and declare the pattern:

public class SplitSentenceBolt extends BaseRichBolt {
    private OutputCollector _collector;
    private Pattern COMMA_PATTERN; // Fixed variable name for clarity

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this._collector = collector;
        this.COMMA_PATTERN = Pattern.compile(",");
    }

    @Override
    public void execute(Tuple input) {
        String sentence = input.getStringByField(RabbitMQSpout.DATA);
        String[] words = COMMA_PATTERN.split(sentence);
        if (words.length > 0) {
            // Trim whitespace from the first field (in case of CSV leading spaces)
            _collector.emit(new Values(words[0].trim()));
        }
        // Ack the tuple to prevent Storm from replaying it
        _collector.ack(input);
    }

    // Rest of your methods remain the same...
}

Also ensure your topology references the correct class name:

builder.setBolt("split", new SplitSentenceBolt(), 1)
       .shuffleGrouping("spout");

4. RabbitMQExporterBolt: Missing Tuple Acknowledgment

Since you're extending BaseRichBolt, you must manually acknowledge tuples to let Storm know processing is complete. Also, the final bolt doesn't need to declare output fields (it's the end of the pipeline):

@Override
public void execute(Tuple tuple) {
    String word = tuple.getString(0);
    Integer count = tuple.getInteger(1);
    String output = word + " " + count;
    rabbitmq.send(output);
    // Critical: Acknowledge the tuple
    collector.ack(tuple);
}

@Override
public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
    // No need to declare fields for a terminal bolt
    outputFieldsDeclarer.declare(new Fields());
}

5. RabbitMQManager: Connection/Channel Cleanup

Your initializeQueue() method creates a new connection/channel every time and closes it immediately — this is inefficient, but more importantly, your send() and receive() methods create a new channel for every operation. While this works, it's better to reuse channels where possible, but the immediate fix is ensuring connections are properly reopened when needed.


Post-Fix Testing Steps

  1. Restart your Storm cluster (local or distributed)
  2. Verify the input queue has messages
  3. Check Storm logs for any remaining errors (enable debug mode in your config to see tuple flow)
  4. Monitor the output queue — you should now see word-count tuples appearing

内容的提问来源于stack exchange,提问作者Simon

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:52:30