如何检测Video是否已在处理,避免重复入队并同步处理结果
实现单视频对应Job的队列控制与结果分发逻辑
你的核心需求是避免同一原始视频重复进入处理队列,同时让所有请求该视频处理结果的调用者能在处理完成后获取到编码后的视频。现有Job类已经通过synchronized和wait/notifyAll实现了多线程等待结果的能力,接下来只需要补充一个全局的Job管理器来跟踪视频的处理状态即可。
实现方案
我们需要一个JobManager类,负责:
- 维护原始视频与对应
Job的映射,确保每个视频唯一对应一个Job - 处理请求时判断视频是否已在处理,返回对应
Job让调用者等待结果 - 管理处理队列,确保未在处理的视频进入队列执行
完整代码实现
import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingQueue; // 假设已存在的Video类 class Video { private String videoId; public Video(String videoId) { this.videoId = videoId; } public String getVideoId() { return videoId; } // 必须重写equals和hashCode,确保同一视频能被Map正确识别 @Override public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; Video video = (Video) o; return videoId.equals(video.videoId); } @Override public int hashCode() { return videoId.hashCode(); } } // 你提供的Job类(保留原有逻辑) public class Job { private Video o_video; private Video coded_video; public Job(Video v) { o_video = v; coded_video = null; } public Video extractOriginal() { return o_video; } public void addCodedVideo(Video vCod) { synchronized(this) { coded_video = vCod; this.notifyAll(); } } public Video extractCodedVideo() throws InterruptedException { synchronized(this) { while(coded_video == null) { this.wait(); } } return coded_video; } } // Job管理器,核心逻辑实现 class JobManager { // 用ConcurrentHashMap保证多线程下的映射安全 private final Map<Video, Job> videoJobMap = new ConcurrentHashMap<>(); // 线程安全的处理队列 private final LinkedBlockingQueue<Job> processingQueue = new LinkedBlockingQueue<>(); public JobManager() { // 启动处理线程 startProcessingWorker(); } // 获取或创建对应视频的Job,处理请求逻辑 public Job getOrCreateJob(Video originalVideo) { // 先尝试从缓存获取 Job existingJob = videoJobMap.get(originalVideo); if (existingJob != null) { return existingJob; } // 双重检查锁定,避免多线程下重复创建Job synchronized (videoJobMap) { existingJob = videoJobMap.get(originalVideo); if (existingJob == null) { existingJob = new Job(originalVideo); videoJobMap.put(originalVideo, existingJob); // 将新Job加入处理队列 processingQueue.offer(existingJob); } } return existingJob; } // 启动处理队列的worker线程 private void startProcessingWorker() { Thread worker = new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { Job job = processingQueue.take(); // 替换为实际的视频编码逻辑 Video original = job.extractOriginal(); Video codedVideo = new Video(original.getVideoId() + "_coded"); // 处理完成后设置编码视频,唤醒所有等待的线程 job.addCodedVideo(codedVideo); // 可选:处理完成后从缓存移除,节省内存(不需要保留历史结果时启用) // videoJobMap.remove(original); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }); worker.setDaemon(true); worker.start(); } } // 测试示例 class Test { public static void main(String[] args) throws InterruptedException { JobManager manager = new JobManager(); Video testVideo = new Video("video_001"); // 模拟5个并发请求获取同一视频的处理结果 for (int i = 0; i < 5; i++) { int requestId = i; new Thread(() -> { try { Job job = manager.getOrCreateJob(testVideo); Video coded = job.extractCodedVideo(); System.out.println("请求" + requestId + "获取到编码视频:" + coded.getVideoId()); } catch (InterruptedException e) { e.printStackTrace(); } }).start(); } Thread.sleep(1000); } }
关键逻辑说明
- 视频唯一性保障:
Video类重写equals()和hashCode(),结合ConcurrentHashMap的映射关系,确保同一视频只会创建一个Job实例。 - 重复请求处理:已有
Job存在时直接返回,调用者调用extractCodedVideo()会进入等待状态,直到处理完成被唤醒。 - 线程安全队列:
LinkedBlockingQueue保证队列操作的线程安全,worker线程从队列取Job执行编码,确保每个Job仅被处理一次。 - 结果分发:
Job类中的notifyAll()唤醒所有等待该结果的线程,实现一次处理、多请求获取结果的效果。
内容的提问来源于stack exchange,提问作者Henry FXP
相关产品推荐
相关产品推荐

