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

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

问题原因及解决思路

核心问题

  1. 字段不匹配:聚合Bolt每次只能接收单个来源的Tuple(要么是高斯模糊帧,要么是锐化帧),原代码直接读取两个来源的字段,导致当Tuple来自高斯模糊Bolt时,sharpenedFrame字段不存在,抛出异常。
  2. 帧配对缺失:Storm的shuffle分组无法保证同编号的两个帧同时到达聚合Bolt,必须手动缓存未配对的帧,等待对应帧到达后再执行聚合。

解决要点

  1. 按来源区分处理Tuple:通过input.getSourceComponent()判断当前Tuple的来源,分别处理两种类型的帧。
  2. 添加缓存机制:用两个HashMap缓存未配对的帧,以帧编号为key,确保同编号帧能配对聚合。
  3. 优化分组策略(可选):将聚合Bolt的shuffle分组改为fieldsGrouping,按帧编号分组,确保同编号帧发送到同一个聚合Bolt实例,避免多实例场景下的跨实例缓存问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:49:52