Apache Storm多Bolt Tuple接收异常:聚合Bolt无法合并双源数据
Apache Storm拓扑帧聚合异常问题解决
我正在构建一个Apache Storm拓扑,用于读取视频样本、格式转换并输出。拓扑逻辑为:Spout读取视频帧后发送至BoltImageProcessor,再分流至BoltGaussianBlur(高斯模糊处理)和BoltSharpener(帧锐化处理)两个Bolt。这两个Bolt均能成功向BoltFrameAggregator发送Tuple,但聚合Bolt每次仅能接收单个源的Tuple,且运行时抛出“sharpenedFrame does not exist”异常,无法完成同编号帧的合并,该如何解决?
Topology.java
import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.topology.TopologyBuilder; import java.io.*; public class Topology { public static void main(String[] args) throws Exception { // 创建日志文件记录标准输出和错误输出 File logFile = new File("topology.log"); // 重定向标准输出和错误输出到日志文件 PrintStream printStream = new PrintStream(new FileOutputStream(logFile)); try (printStream) { System.setOut(printStream); System.setErr(printStream); Config config = new Config(); // 创建集群配置实例 config.setDebug(true); TopologyBuilder builder = new TopologyBuilder(); // 创建TopologyBuilder Spout spout = new Spout(); // 初始化Spout和Bolts BoltFrameAnalyzer boltFrameAnalyzer = new BoltFrameAnalyzer(); BoltAnalysisSaver boltAnalysisSaver = new BoltAnalysisSaver(); BoltImageProcessor boltImageProcessor = new BoltImageProcessor(); BoltGaussianBlur boltGaussianBlur = new BoltGaussianBlur(); BoltSharpener boltSharpener = new BoltSharpener(); BoltFrameAggregator boltFrameAggregator = new BoltFrameAggregator(); BoltOutputGenerator boltOutputGenerator = new BoltOutputGenerator(); builder.setSpout("spout", spout, 1); // 定义数据流连接关系 builder.setBolt("bolt-frame-analyzer", boltFrameAnalyzer, 1).shuffleGrouping("spout"); builder.setBolt("bolt-analysis-saver", boltAnalysisSaver, 1).shuffleGrouping("bolt-frame-analyzer"); builder.setBolt("bolt-image-processor", boltImageProcessor, 1).shuffleGrouping("spout"); builder.setBolt("bolt-gaussian-blur", boltGaussianBlur, 1).shuffleGrouping("bolt-image-processor"); builder.setBolt("bolt-sharpener", boltSharpener, 1).shuffleGrouping("bolt-image-processor"); // 可选优化:改为按帧编号分组,确保同编号帧到同一聚合实例 builder.setBolt("bolt-frame-aggregator", boltFrameAggregator, 1) .fieldsGrouping("bolt-sharpener", new Fields("sharpenedFrameNumber")) .fieldsGrouping("bolt-gaussian-blur", new Fields("gaussianBlurFrameNumber")); builder.setBolt("bolt-output-generator", boltOutputGenerator, 1).shuffleGrouping("bolt-frame-aggregator"); try (LocalCluster cluster = new LocalCluster()) { // 使用try-with-resources自动关闭集群 cluster.submitTopology("Topology", config, builder.createTopology()); Thread.sleep(100000); // 根据需要调整休眠时间 } // 退出try块时自动关闭集群 } // 关闭PrintStream和日志文件 } }
BoltGaussianBlur.java
import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import org.opencv.core.Mat; import org.opencv.core.Size; import org.opencv.imgcodecs.Imgcodecs; import org.opencv.imgproc.Imgproc; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.FileInputStream; import java.io.IOException; import java.util.Properties; public class BoltGaussianBlur extends BaseBasicBolt { private final String framesGaussianBlurFilePath; public BoltGaussianBlur() { try { Properties prop = new Properties(); prop.load(new FileInputStream("config.ini")); framesGaussianBlurFilePath = prop.getProperty("framesGaussianBlurFilePath"); } catch (IOException e) { LOG.error("BoltImageProcessor: 读取config.ini文件时出错", e); throw new RuntimeException(e); } } private static final Logger LOG = LoggerFactory.getLogger(BoltGaussianBlur.class); public void execute(Tuple input, BasicOutputCollector collector) { int gaussianBlurFrameNumber = input.getIntegerByField("frameNumber"); LOG.info("BoltGaussianBlur: 成功接收帧 #" + gaussianBlurFrameNumber); Mat receivedFrame = (Mat) input.getValueByField("resizedFrame"); Mat gaussianBlurFrame = new Mat(); Imgproc.GaussianBlur(receivedFrame, gaussianBlurFrame, new Size(9, 9), 2, 2); String gaussianBlurFileName = framesGaussianBlurFilePath + "/frame_" + gaussianBlurFrameNumber + "_gaussian_blur.png"; Imgcodecs.imwrite(gaussianBlurFileName, gaussianBlurFrame); LOG.info("BoltGaussianBlur: 帧 #" + gaussianBlurFrameNumber + " 已成功转换为高斯模糊格式"); collector.emit(new Values(gaussianBlurFrame, gaussianBlurFrameNumber)); } public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("gaussianBlurFrame", "gaussianBlurFrameNumber")); } }
BoltSharpener.java
import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import org.opencv.core.Core; import org.opencv.core.Mat; import org.opencv.core.Size; import org.opencv.imgproc.Imgproc; import org.opencv.imgcodecs.Imgcodecs; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.FileInputStream; import java.io.IOException; import java.util.Properties; public class BoltSharpener extends BaseBasicBolt { private final String framesSharpenedFilePath; public BoltSharpener() { try { Properties prop = new Properties(); prop.load(new FileInputStream("config.ini")); framesSharpenedFilePath = prop.getProperty("framesSharpenedFilePath"); } catch (IOException e) { LOG.error("BoltSharpener: 读取config.ini文件时出错", e); throw new RuntimeException(e); } } private static final Logger LOG = LoggerFactory.getLogger(BoltSharpener.class); public void execute(Tuple input, BasicOutputCollector collector) { int sharpenedFrameNumber = input.getIntegerByField("frameNumber"); LOG.info("BoltSharpener: 成功接收帧 #" + sharpenedFrameNumber); Mat receivedFrame = (Mat) input.getValueByField("resizedFrame"); Mat sharpenedFrame = new Mat(); Imgproc.GaussianBlur(receivedFrame, sharpenedFrame, new Size(0, 0), 10); Core.addWeighted(receivedFrame, 1.5, sharpenedFrame, -0.5, 0, sharpenedFrame); String sharpenedFileName = framesSharpenedFilePath + "/frame_" + sharpenedFrameNumber + "_sharpened.png"; Imgcodecs.imwrite(sharpenedFileName, sharpenedFrame); LOG.info("BoltSharpener: 帧 #" + sharpenedFrameNumber + " 已成功锐化"); collector.emit(new Values(sharpenedFrame, sharpenedFrameNumber)); } public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("sharpenedFrame", "sharpenedFrameNumber")); } }
修改后的BoltFrameAggregator.java
import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import org.opencv.core.Core; import org.opencv.core.Mat; import org.opencv.imgcodecs.Imgcodecs; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.FileInputStream; import java.io.IOException; import java.util.HashMap; import java.util.Map; import java.util.Properties; public class BoltFrameAggregator extends BaseBasicBolt { private final String framesAggregated; // 缓存未配对的帧,key为帧编号 private final Map<Integer, Mat> gaussianBlurCache = new HashMap<>(); private final Map<Integer, Mat> sharpenedCache = new HashMap<>(); public BoltFrameAggregator() { try { Properties prop = new Properties(); prop.load(new FileInputStream("config.ini")); framesAggregated = prop.getProperty("framesAggregatedFilePath"); } catch (IOException e) { LOG.error("BoltAggregator: 读取config.ini文件时出错", e); throw new RuntimeException(e); } } private static final Logger LOG = LoggerFactory.getLogger(BoltFrameAggregator.class); public void execute(Tuple input, BasicOutputCollector collector) { // 根据来源区分处理不同类型的帧 if (input.getSourceComponent().equals("bolt-gaussian-blur")) { Mat gaussianFrame = (Mat) input.getValueByField("gaussianBlurFrame"); int frameNum = input.getIntegerByField("gaussianBlurFrameNumber"); // 检查是否有对应编号的锐化帧 if (sharpenedCache.containsKey(frameNum)) { Mat sharpenedFrame = sharpenedCache.remove(frameNum); aggregateFrames(sharpenedFrame, gaussianFrame, frameNum, collector); } else { // 无配对帧,存入缓存 gaussianBlurCache.put(frameNum, gaussianFrame); LOG.info("BoltFrameAggregator: 缓存高斯模糊帧 #" + frameNum); } } else if (input.getSourceComponent().equals("bolt-sharpener")) { Mat sharpenedFrame = (Mat) input.getValueByField("sharpenedFrame"); int frameNum = input.getIntegerByField("sharpenedFrameNumber"); // 检查是否有对应编号的高斯模糊帧 if (gaussianBlurCache.containsKey(frameNum)) { Mat gaussianFrame = gaussianBlurCache.remove(frameNum); aggregateFrames(sharpenedFrame, gaussianFrame, frameNum, collector); } else { // 无配对帧,存入缓存 sharpenedCache.put(frameNum, sharpenedFrame); LOG.info("BoltFrameAggregator: 缓存锐化帧 #" + frameNum); } } } // 单独提取聚合逻辑,提升可读性 private void aggregateFrames(Mat sharpenedFrame, Mat gaussianFrame, int frameNum, BasicOutputCollector collector) { Mat aggregatedFrame = new Mat(); Core.addWeighted(sharpenedFrame, 1, gaussianFrame, 1, 0, aggregatedFrame); String aggregatedFileName = framesAggregated + "/frame_" + frameNum + "_aggregated.png"; Imgcodecs.imwrite(aggregatedFileName, aggregatedFrame); LOG.info("BoltFrameAggregator: 帧 #" + frameNum + " 已成功聚合"); collector.emit(new Values(aggregatedFrame, frameNum)); } public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("aggregatedFrame", "aggregatedFrameNumber")); } }
原始错误日志
[Thread-42-bolt-frame-aggregator-executor[3, 3]] INFO o.a.s.e.Executor - Processing received TUPLE: source: bolt-gaussian-blur:5, stream: default, id: {}, [Mat [ 720*1280*CV_8UC1, isCont=true, isSubmat=false, nativeObj=0x1fc231bd4a0, dataAddr=0x1fc258f0f60 ], 0] PROC_START_TIME(sampled): null EXEC_START_TIME(sampled): null for TASK: 3 ... [Thread-42-bolt-frame-aggregator-executor[3, 3]] ERROR o.a.s.u.Utils - Async loop died! java.lang.RuntimeException: java.lang.IllegalArgumentException: sharpenedFrame does not exist at org.apache.storm.executor.Executor.accept(Executor.java:301) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.utils.JCQueue.consumeImpl(JCQueue.java:113) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.utils.JCQueue.consume(JCQueue.java:89) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.executor.bolt.BoltExecutor$1.call(BoltExecutor.java:154) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.executor.bolt.BoltExecutor$1.call(BoltExecutor.java:140) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.utils.Utils$1.run(Utils.java:398) ~[storm-client-2.6.0.jar:2.6.0] at java.base/java.lang.Thread.run(Thread.java:1583) [?:?] Caused by: java.lang.IllegalArgumentException: sharpenedFrame does not exist at org.apache.storm.tuple.Fields.fieldIndex(Fields.java:98) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.tuple.TupleImpl.fieldIndex(TupleImpl.java:101) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.tuple.TupleImpl.getValueByField(TupleImpl.java:161) ~[storm-client-2.6.0.jar:2.6.0] at BoltFrameAggregator.execute(BoltFrameAggregator.java:33) ~[classes/:?] at org.apache.storm.topology.BasicBoltExecutor.execute(BasicBoltExecutor.java:48) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.executor.bolt.BoltExecutor.tupleActionFn(BoltExecutor.java:212) ~[storm-client-2.6.0.jar:2.6.0] at org.apache.storm.executor.Executor.accept(Executor.java:294) ~[storm-client-2.6.0.jar:2.6.0] ... 6 more
问题原因及解决思路
核心问题
- 字段不匹配:聚合Bolt每次只能接收单个来源的Tuple(要么是高斯模糊帧,要么是锐化帧),原代码直接读取两个来源的字段,导致当Tuple来自高斯模糊Bolt时,
sharpenedFrame字段不存在,抛出异常。 - 帧配对缺失:Storm的shuffle分组无法保证同编号的两个帧同时到达聚合Bolt,必须手动缓存未配对的帧,等待对应帧到达后再执行聚合。
解决要点
- 按来源区分处理Tuple:通过
input.getSourceComponent()判断当前Tuple的来源,分别处理两种类型的帧。 - 添加缓存机制:用两个HashMap缓存未配对的帧,以帧编号为key,确保同编号帧能配对聚合。
- 优化分组策略(可选):将聚合Bolt的shuffle分组改为
fieldsGrouping,按帧编号分组,确保同编号帧发送到同一个聚合Bolt实例,避免多实例场景下的跨实例缓存问题。
内容的提问来源于stack exchange,提问作者Hatam Abolghasemi
相关产品推荐
相关产品推荐

