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

Java三线程通信咨询:音频下载/播放与WebSocket通知协作

Java三线程协作通信实现方案

这是一个典型的多线程协作场景,结合生产者-消费者模式和事件触发机制就能很好地解决三者间的通信问题。我给你梳理一下具体的实现思路和可运行的代码示例,你可以根据实际业务需求调整细节:

核心通信逻辑梳理

  1. ScheduleDownloader → AudioPlayer:用线程安全的阻塞队列实现生产者-消费者模型,下载完成后直接将音频调度对象放入队列,AudioPlayer会自动阻塞等待新内容,天然完成“下载完成通知播放”的需求。
  2. WebSocketListener → ScheduleDownloader:用Condition信号实现触发通知,WebSocket监听到刷新消息后,给下载线程发送唤醒信号,触发重新下载流程。
  3. 整体协同:三个线程共享线程安全的状态变量,确保不会出现竞态条件,同时每个线程维护自己的生命周期,支持优雅停止。

具体代码实现

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. 阻塞队列的容量控制:设置容量为1,确保当有新的音频调度时,旧的过时内容会被替换,避免AudioPlayer播放过时内容。
  2. Condition替代轮询:用Lock+Condition实现高效的等待/唤醒,比sleep轮询更节省资源,响应更及时。
  3. 优雅停止机制:每个线程都通过volatile标记控制运行状态,同时处理中断信号,确保线程能安全退出。
  4. 播放中断支持:AudioPlayer中用volatile标记控制播放状态,收到新调度时能立即停止当前播放,切换到最新内容。

内容的提问来源于stack exchange,提问作者Daniel Haughton

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:45:52