Java三线程通信咨询:音频下载/播放与WebSocket通知协作
Java三线程协作通信实现方案
这是一个典型的多线程协作场景,结合生产者-消费者模式和事件触发机制就能很好地解决三者间的通信问题。我给你梳理一下具体的实现思路和可运行的代码示例,你可以根据实际业务需求调整细节:
核心通信逻辑梳理
- ScheduleDownloader → AudioPlayer:用线程安全的阻塞队列实现生产者-消费者模型,下载完成后直接将音频调度对象放入队列,AudioPlayer会自动阻塞等待新内容,天然完成“下载完成通知播放”的需求。
- WebSocketListener → ScheduleDownloader:用Condition信号实现触发通知,WebSocket监听到刷新消息后,给下载线程发送唤醒信号,触发重新下载流程。
- 整体协同:三个线程共享线程安全的状态变量,确保不会出现竞态条件,同时每个线程维护自己的生命周期,支持优雅停止。
具体代码实现
1. 定义核心数据类
首先定义封装音频和调度信息的实体类:
// 音频调度配置类,自定义播放规则、时间等 class ScheduleConfig { // 这里可以添加你的调度字段,比如播放时长、循环次数、播放顺序等 } // 音频+调度的组合对象 class AudioSchedule { private byte[] audioData; // 下载的音频二进制数据 private ScheduleConfig scheduleConfig; // 对应的调度配置 public AudioSchedule(byte[] audioData, ScheduleConfig scheduleConfig) { this.audioData = audioData; this.scheduleConfig = scheduleConfig; } // Getter方法,供AudioPlayer使用 public byte[] getAudioData() { return audioData; } public ScheduleConfig getScheduleConfig() { return scheduleConfig; } }
2. ScheduleDownloader线程(下载+通知播放)
负责下载内容,并通过阻塞队列通知AudioPlayer:
import java.util.concurrent.BlockingQueue; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; public class ScheduleDownloader extends Thread { private final BlockingQueue<AudioSchedule> audioQueue; private volatile boolean running = true; private boolean needRefresh = false; private final Lock lock = new ReentrantLock(); private final Condition refreshCondition = lock.newCondition(); public ScheduleDownloader(BlockingQueue<AudioSchedule> audioQueue) { this.audioQueue = audioQueue; } @Override public void run() { while (running) { try { // 1. 下载音频和调度对象(模拟实际下载逻辑) AudioSchedule newSchedule = downloadAudioAndSchedule(); // 2. 将新内容放入队列,若队列已满则替换旧内容(确保只保留最新调度) if (!audioQueue.offer(newSchedule)) { audioQueue.poll(); // 移除旧的过时内容 audioQueue.put(newSchedule); } System.out.println("[ScheduleDownloader] 新音频调度已下载并提交"); // 3. 等待刷新信号或停止指令(用Condition替代轮询sleep,更高效) lock.lock(); try { while (running && !needRefresh) { refreshCondition.await(); } needRefresh = false; // 重置刷新标记 } finally { lock.unlock(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[ScheduleDownloader] 线程被中断"); break; } catch (Exception e) { System.err.println("[ScheduleDownloader] 下载失败,5秒后重试:" + e.getMessage()); try { Thread.sleep(5000); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } } } // WebSocket线程调用此方法触发重新下载 public void triggerRefresh() { lock.lock(); try { needRefresh = true; refreshCondition.signal(); // 唤醒等待的下载线程 } finally { lock.unlock(); } } // 模拟实际下载逻辑(替换为你的HTTP/网络下载代码) private AudioSchedule downloadAudioAndSchedule() throws InterruptedException { System.out.println("[ScheduleDownloader] 开始下载音频和调度..."); Thread.sleep(2000); // 模拟下载耗时 return new AudioSchedule(new byte[1024], new ScheduleConfig()); } // 优雅停止线程 public void stopGracefully() { running = false; triggerRefresh(); // 唤醒等待的线程,确保能及时退出 this.interrupt(); } }
3. AudioPlayer线程(接收通知+播放音频)
阻塞等待新的音频调度,收到后立即播放(支持中断当前播放切换新内容):
import java.util.concurrent.BlockingQueue; public class AudioPlayer extends Thread { private final BlockingQueue<AudioSchedule> audioQueue; private volatile boolean running = true; private volatile boolean isPlaying = false; // 标记当前是否在播放 public AudioPlayer(BlockingQueue<AudioSchedule> audioQueue) { this.audioQueue = audioQueue; } @Override public void run() { while (running) { try { // 阻塞等待新的音频调度 AudioSchedule newSchedule = audioQueue.take(); System.out.println("[AudioPlayer] 收到新音频调度,停止当前播放"); isPlaying = false; // 中断当前播放任务 // 开始播放新内容 System.out.println("[AudioPlayer] 开始播放新音频"); playAudio(newSchedule); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[AudioPlayer] 线程被中断"); break; } catch (Exception e) { System.err.println("[AudioPlayer] 播放失败,等待下一个调度:" + e.getMessage()); } } } // 模拟音频播放逻辑(替换为你的音频播放库调用) private void playAudio(AudioSchedule schedule) throws InterruptedException { isPlaying = true; try { // 模拟分块播放,支持中途停止 for (int i = 0; i < 5; i++) { if (!isPlaying) break; Thread.sleep(1000); System.out.println("[AudioPlayer] 播放中..."); } } finally { isPlaying = false; } System.out.println("[AudioPlayer] 当前音频播放完成"); } // 优雅停止线程 public void stopGracefully() { running = false; isPlaying = false; // 停止当前播放 this.interrupt(); } }
4. WebSocketListener线程(监听刷新消息+触发下载)
监听WebSocket消息,收到“内容过时”通知后触发ScheduleDownloader重新下载:
public class WebSocketListener extends Thread { private final ScheduleDownloader downloader; private volatile boolean running = true; public WebSocketListener(ScheduleDownloader downloader) { this.downloader = downloader; } @Override public void run() { while (running) { try { // 模拟监听WebSocket消息(替换为你的WebSocket客户端逻辑) String message = listenWebSocketMessage(); if ("CONTENT_OUTDATED".equals(message)) { // 假设刷新消息标识为这个 System.out.println("[WebSocketListener] 收到内容过时通知,触发重新下载"); downloader.triggerRefresh(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("[WebSocketListener] 线程被中断"); break; } catch (Exception e) { System.err.println("[WebSocketListener] 监听出错,3秒后重试:" + e.getMessage()); try { Thread.sleep(3000); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } } } // 模拟WebSocket消息监听 private String listenWebSocketMessage() throws InterruptedException { System.out.println("[WebSocketListener] 等待WebSocket消息..."); Thread.sleep(10000); // 模拟等待消息的间隔 return "CONTENT_OUTDATED"; // 模拟收到刷新通知 } // 优雅停止线程 public void stopGracefully() { running = false; this.interrupt(); } }
5. 应用启动类
整合三个线程,启动并管理生命周期:
import java.util.concurrent.LinkedBlockingQueue; public class AudioApplication { public static void main(String[] args) throws InterruptedException { // 初始化阻塞队列,容量设为1确保只保留最新的音频调度 LinkedBlockingQueue<AudioSchedule> audioQueue = new LinkedBlockingQueue<>(1); // 创建线程实例 ScheduleDownloader downloader = new ScheduleDownloader(audioQueue); AudioPlayer player = new AudioPlayer(audioQueue); WebSocketListener listener = new WebSocketListener(downloader); // 启动线程 downloader.start(); player.start(); listener.start(); // 模拟应用运行30秒后停止 Thread.sleep(30000); System.out.println("\n[应用] 准备优雅停止..."); downloader.stopGracefully(); player.stopGracefully(); listener.stopGracefully(); // 等待所有线程结束 downloader.join(); player.join(); listener.join(); System.out.println("[应用] 已停止"); } }
关键优化点说明
- 阻塞队列的容量控制:设置容量为1,确保当有新的音频调度时,旧的过时内容会被替换,避免AudioPlayer播放过时内容。
- Condition替代轮询:用Lock+Condition实现高效的等待/唤醒,比sleep轮询更节省资源,响应更及时。
- 优雅停止机制:每个线程都通过volatile标记控制运行状态,同时处理中断信号,确保线程能安全退出。
- 播放中断支持:AudioPlayer中用volatile标记控制播放状态,收到新调度时能立即停止当前播放,切换到最新内容。
内容的提问来源于stack exchange,提问作者Daniel Haughton
相关产品推荐
相关产品推荐

