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

Java单进程下Player实例消息收发实现技术求助

实现思路与最小可运行代码

没问题,我来帮你理清楚这个需求的实现思路和代码~

首先,咱们的核心需求是同一进程内的两个Player线程互相通信,用生产者消费者模式完全可行,关键是用线程安全的阻塞队列作为它们的消息通道——因为阻塞队列自带线程安全特性,而且当队列空的时候,取消息的线程会自动阻塞,不用自己写等待逻辑,非常适合这个场景。

具体的设计逻辑:

  • 每个Player实例拥有自己的消息接收队列(用来接收对方发的消息),同时持有对另一个Player的接收队列引用(用来发送消息给对方)。
  • 每个Player维护一个发送计数器sendCounter,每发送一条消息就自增。
  • 每个Player启动一个独立线程,循环监听自己的接收队列:一旦收到消息,就按照要求拼接回复内容(原消息 + 自己当前的发送计数器值),然后发送给对方。

接下来是最小可运行的纯Java代码:

代码实现

Player类

import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

public class Player implements Runnable {
    // 自身的消息接收队列
    private final BlockingQueue<String> receiveQueue;
    // 目标Player的消息接收队列(用来发送消息给对方)
    private final BlockingQueue<String> targetReceiveQueue;
    // 发送消息计数器
    private int sendCounter = 0;
    // Player名称,方便日志区分
    private final String name;

    public Player(String name, BlockingQueue<String> targetReceiveQueue) {
        this.name = name;
        this.receiveQueue = new LinkedBlockingQueue<>();
        this.targetReceiveQueue = targetReceiveQueue;
        // 启动监听线程
        new Thread(this).start();
    }

    // 发送消息方法
    public void send(String message) {
        try {
            sendCounter++;
            System.out.printf("[%s] 发送消息: %s (当前发送计数器: %d)%n", name, message, sendCounter);
            targetReceiveQueue.put(message);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    // 线程核心逻辑:循环接收并处理消息
    @Override
    public void run() {
        while (!Thread.currentThread().isInterrupted()) {
            try {
                // 阻塞等待接收消息
                String receivedMsg = receiveQueue.take();
                System.out.printf("[%s] 收到消息: %s%n", name, receivedMsg);
                // 拼接回复内容:原消息 + 自身发送计数器值
                String replyMsg = receivedMsg + " | 我的发送计数器: " + sendCounter;
                // 发送回复
                send(replyMsg);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
    }

    // 获取自身的接收队列(用于给另一个Player绑定目标队列)
    public BlockingQueue<String> getReceiveQueue() {
        return receiveQueue;
    }
}

主类(测试入口)

public class PlayerCommunicationDemo {
    public static void main(String[] args) throws InterruptedException {
        // 创建两个Player的接收队列(先初始化队列,再互相绑定)
        BlockingQueue<String> player1Queue = new LinkedBlockingQueue<>();
        BlockingQueue<String> player2Queue = new LinkedBlockingQueue<>();

        // 初始化两个Player,互相指定对方的接收队列作为发送目标
        Player initiator = new Player("发起方Player", player2Queue);
        Player responder = new Player("响应方Player", player1Queue);

        // 发起方发送第一条消息
        initiator.send("这是第一条测试消息");

        // 让程序运行一段时间,观察交互
        Thread.sleep(2000);

        // 中断线程,结束程序
        Thread.currentThread().interrupt();
    }
}

代码说明

  1. 阻塞队列的作用:LinkedBlockingQueue是线程安全的阻塞队列,take()方法会在队列空时阻塞线程,put()方法会在队列满时阻塞(这里我们用无界队列,所以不会满),完美解决了线程间的同步问题。
  2. Player的线程逻辑:每个Player启动后就进入循环监听自己的接收队列,收到消息就处理并回复,完全符合“收到消息就回复”的需求。
  3. 计数器的维护:每次调用send()方法时,计数器自增,回复时把当前的计数器值拼接到消息里,满足需求。

运行这个代码,你会看到类似这样的输出:

[发起方Player] 发送消息: 这是第一条测试消息 (当前发送计数器: 1)
[响应方Player] 收到消息: 这是第一条测试消息
[响应方Player] 发送消息: 这是第一条测试消息 | 我的发送计数器: 1 (当前发送计数器: 1)
[发起方Player] 收到消息: 这是第一条测试消息 | 我的发送计数器: 1
[发起方Player] 发送消息: 这是第一条测试消息 | 我的发送计数器: 1 | 我的发送计数器: 2 (当前发送计数器: 2)
[响应方Player] 收到消息: 这是第一条测试消息 | 我的发送计数器: 1 | 我的发送计数器: 2
...

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 15:57:57