基于Storm与RabbitMQ的单词计数拓扑无输出问题求助
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
- Restart your Storm cluster (local or distributed)
- Verify the
inputqueue has messages - Check Storm logs for any remaining errors (enable debug mode in your config to see tuple flow)
- Monitor the
outputqueue — you should now see word-count tuples appearing
内容的提问来源于stack exchange,提问作者Simon

