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

如何基于Java与gRPC实现双客户端聊天及消息互传功能?

基于Java实现gRPC双向聊天功能的方案

一、定义gRPC协议(.proto文件)

首先编写Protocol Buffers定义,明确服务接口与消息类型,这是gRPC通信的核心基础:

syntax = "proto3";

package chat;

option java_package = "com.example.chat.grpc";
option java_outer_classname = "ChatProto";
option java_multiple_files = true;

// 聊天服务接口定义
service ChatService {
    // 客户端登录,同时建立双向流接收服务器推送
    rpc Login(LoginRequest) returns (stream ServerNotification);
    // 请求创建双人聊天室
    rpc CreateRoom(RoomCreateRequest) returns (EmptyResponse);
    // 发送聊天消息
    rpc SendMessage(ChatMessage) returns (EmptyResponse);
}

// 登录请求
message LoginRequest {
    string user_id = 1; // 客户端唯一标识(如用户名)
}

// 创建房间请求
message RoomCreateRequest {
    string requester_id = 1; // 请求发起者ID
    string target_user_id = 2; // 目标用户ID
}

// 服务器推送的通知内容(包含房间配置或聊天消息)
message ServerNotification {
    oneof content {
        RoomConfig room_config = 1;
        ChatMessage chat_message = 2;
    }
}

// 聊天室配置信息
message RoomConfig {
    string room_id = 1;
    string peer_user_id = 2; // 聊天对方的用户ID
}

// 聊天消息体
message ChatMessage {
    string room_id = 1;
    string sender_id = 2;
    string content = 3;
}

// 空响应
message EmptyResponse {}

二、生成Java代码

使用protobuf编译器生成对应Java代码,执行以下命令(需确保protobuf及gRPC插件已配置):

protoc --java_out=. --grpc-java_out=. --plugin=protoc-gen-grpc-java=/path/to/protoc-gen-grpc-java chat.proto

三、服务器端实现

服务器需要维护客户端连接状态,处理登录、房间创建与消息转发逻辑:

1. 客户端连接管理类

import com.example.chat.grpc.ChatProto.*;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.ConcurrentHashMap;

public class ClientManager {
    private static final ConcurrentHashMap<String, StreamObserver<ServerNotification>> clientStubs = new ConcurrentHashMap<>();

    // 注册客户端的通知接收流
    public static void registerClient(String userId, StreamObserver<ServerNotification> observer) {
        clientStubs.put(userId, observer);
    }

    // 获取指定用户的通知接收流
    public static StreamObserver<ServerNotification> getClientStub(String userId) {
        return clientStubs.get(userId);
    }

    // 移除断开连接的客户端
    public static void removeClient(String userId) {
        clientStubs.remove(userId);
    }
}

2. gRPC服务实现类

import com.example.chat.grpc.ChatProto.*;
import com.example.chat.grpc.ChatServiceGrpc;
import io.grpc.stub.StreamObserver;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;

public class ChatServiceImpl extends ChatServiceGrpc.ChatServiceImplBase {
    // 维护房间与成员的映射,用于快速查找聊天对象
    private final ConcurrentHashMap<String, String[]> roomMembers = new ConcurrentHashMap<>();

    @Override
    public void login(LoginRequest request, StreamObserver<ServerNotification> responseObserver) {
        String userId = request.getUserId();
        ClientManager.registerClient(userId, responseObserver);
        // 监听客户端断开事件,清理连接记录
        responseObserver.onCompleted();
        ClientManager.removeClient(userId);
    }

    @Override
    public void createRoom(RoomCreateRequest request, StreamObserver<EmptyResponse> responseObserver) {
        String requesterId = request.getRequesterId();
        String targetId = request.getTargetUserId();

        StreamObserver<ServerNotification> requesterObserver = ClientManager.getClientStub(requesterId);
        StreamObserver<ServerNotification> targetObserver = ClientManager.getClientStub(targetId);

        if (requesterObserver == null || targetObserver == null) {
            responseObserver.onError(new RuntimeException("目标用户未在线"));
            return;
        }

        // 生成唯一房间ID
        String roomId = UUID.randomUUID().toString();
        // 记录房间成员
        roomMembers.put(roomId, new String[]{requesterId, targetId});

        // 向请求者发送房间配置
        RoomConfig requesterConfig = RoomConfig.newBuilder()
                .setRoomId(roomId)
                .setPeerUserId(targetId)
                .build();
        requesterObserver.onNext(ServerNotification.newBuilder().setRoomConfig(requesterConfig).build());

        // 向目标用户发送房间配置
        RoomConfig targetConfig = RoomConfig.newBuilder()
                .setRoomId(roomId)
                .setPeerUserId(requesterId)
                .build();
        targetObserver.onNext(ServerNotification.newBuilder().setRoomConfig(targetConfig).build());

        responseObserver.onNext(EmptyResponse.newBuilder().build());
        responseObserver.onCompleted();
    }

    @Override
    public void sendMessage(ChatMessage request, StreamObserver<EmptyResponse> responseObserver) {
        String roomId = request.getRoomId();
        String senderId = request.getSenderId();
        String[] members = roomMembers.get(roomId);

        if (members == null) {
            responseObserver.onError(new RuntimeException("房间不存在"));
            return;
        }

        // 确定聊天对象ID
        String peerId = members[0].equals(senderId) ? members[1] : members[0];
        StreamObserver<ServerNotification> peerObserver = ClientManager.getClientStub(peerId);

        if (peerObserver != null) {
            ChatMessage message = ChatMessage.newBuilder()
                    .setRoomId(roomId)
                    .setSenderId(senderId)
                    .setContent(request.getContent())
                    .build();
            peerObserver.onNext(ServerNotification.newBuilder().setChatMessage(message).build());
        }

        responseObserver.onNext(EmptyResponse.newBuilder().build());
        responseObserver.onCompleted();
    }
}

3. 启动服务器

import io.grpc.Server;
import io.grpc.ServerBuilder;
import java.io.IOException;

public class ChatServer {
    private Server server;

    private void start() throws IOException {
        int port = 50051;
        server = ServerBuilder.forPort(port)
                .addService(new ChatServiceImpl())
                .build()
                .start();
        System.out.println("服务器启动,监听端口:" + port);
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            System.err.println("*** JVM即将关闭,正在停止gRPC服务器");
            ChatServer.this.stop();
            System.err.println("*** 服务器已关闭");
        }));
    }

    private void stop() {
        if (server != null) {
            server.shutdown();
        }
    }

    private void blockUntilShutdown() throws InterruptedException {
        if (server != null) {
            server.awaitTermination();
        }
    }

    public static void main(String[] args) throws IOException, InterruptedException {
        final ChatServer server = new ChatServer();
        server.start();
        server.blockUntilShutdown();
    }
}

四、客户端实现

客户端需要处理登录、接收服务器推送、发送消息等交互逻辑:

1. 服务器通知处理器

import com.example.chat.grpc.ChatProto.*;
import io.grpc.stub.StreamObserver;

public class ClientNotificationHandler implements StreamObserver<ServerNotification> {
    private final String userId;

    public ClientNotificationHandler(String userId) {
        this.userId = userId;
    }

    @Override
    public void onNext(ServerNotification notification) {
        if (notification.hasRoomConfig()) {
            RoomConfig config = notification.getRoomConfig();
            System.out.println("[" + userId + "] 已加入房间:" + config.getRoomId() + ",聊天对象:" + config.getPeerUserId());
        } else if (notification.hasChatMessage()) {
            ChatMessage message = notification.getChatMessage();
            System.out.println("[" + message.getSenderId() + "]:" + message.getContent());
        }
    }

    @Override
    public void onError(Throwable t) {
        System.err.println("[" + userId + "] 接收通知出错:" + t.getMessage());
    }

    @Override
    public void onCompleted() {
        System.out.println("[" + userId + "] 与服务器连接已断开");
    }
}

2. 客户端主类

import com.example.chat.grpc.ChatProto.*;
import com.example.chat.grpc.ChatServiceGrpc;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import java.util.Scanner;

public class ChatClient {
    private final ManagedChannel channel;
    private final ChatServiceGrpc.ChatServiceStub asyncStub;
    private final String userId;
    private String currentRoomId; // 记录当前所在房间ID

    public ChatClient(String host, int port, String userId) {
        this.channel = ManagedChannelBuilder.forAddress(host, port)
                .usePlaintext()
                .build();
        this.asyncStub = ChatServiceGrpc.newStub(channel);
        this.userId = userId;
    }

    // 登录并启动通知监听
    public void login() {
        LoginRequest request = LoginRequest.newBuilder().setUserId(userId).build();
        ClientNotificationHandler handler = new ClientNotificationHandler(userId);
        asyncStub.login(request, handler);
        System.out.println("[" + userId + "] 已登录服务器");
    }

    // 请求创建房间
    public void createRoom(String targetUserId) {
        RoomCreateRequest request = RoomCreateRequest.newBuilder()
                .setRequesterId(userId)
                .setTargetUserId(targetUserId)
                .build();
        asyncStub.createRoom(request, new StreamObserver<EmptyResponse>() {
            @Override
            public void onNext(EmptyResponse emptyResponse) {
                System.out.println("[" + userId + "] 已发送创建房间请求");
            }

            @Override
            public void onError(Throwable t) {
                System.err.println("[" + userId + "] 创建房间失败:" + t.getMessage());
            }

            @Override
            public void onCompleted() {
            }
        });
    }

    // 发送聊天消息
    public void sendMessage(String content) {
        if (currentRoomId == null) {
            System.err.println("[" + userId + "] 请先加入房间");
            return;
        }
        ChatMessage request = ChatMessage.newBuilder()
                .setRoomId(currentRoomId)
                .setSenderId(userId)
                .setContent(content)
                .build();
        asyncStub.sendMessage(request, new StreamObserver<EmptyResponse>() {
            @Override
            public void onNext(EmptyResponse emptyResponse) {
                System.out.println("[" + userId + "] 已发送:" + content);
            }

            @Override
            public void onError(Throwable t) {
                System.err.println("[" + userId + "] 发送消息失败:" + t.getMessage());
            }

            @Override
            public void onCompleted() {
            }
        });
    }

    public void setCurrentRoomId(String currentRoomId) {
        this.currentRoomId = currentRoomId;
    }

    public void shutdown() throws InterruptedException {
        channel.shutdown().awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS);
    }

    public static void main(String[] args) throws InterruptedException {
        Scanner scanner = new Scanner(System.in);
        System.out.print("输入你的用户ID:");
        String userId = scanner.nextLine();
        ChatClient client = new ChatClient("localhost", 50051, userId);
        client.login();

        // 简单交互逻辑
        while (true) {
            System.out.println("输入指令:1-创建房间 2-发送消息 3-退出");
            String cmd = scanner.nextLine();
            switch (cmd) {
                case "1":
                    System.out.print("输入目标用户ID:");
                    String targetId = scanner.nextLine();
                    client.createRoom(targetId);
                    break;
                case "2":
                    System.out.print("输入消息内容:");
                    String content = scanner.nextLine();
                    client.sendMessage(content);
                    break;
                case "3":
                    client.shutdown();
                    return;
                default:
                    System.out.println("无效指令,请重新输入");
            }
        }
    }
}

五、测试流程

  1. 启动ChatServer;
  2. 启动第一个客户端,输入用户ID(如client1)完成登录;
  3. 启动第二个客户端,输入用户ID(如client2)完成登录;
  4. 在client1控制台输入指令1,输入目标用户IDclient2发起房间创建请求;
  5. 两个客户端会收到服务器推送的房间配置信息,自动记录房间ID;
  6. 在client1输入指令2,输入消息内容,client2会实时收到该消息;
  7. client2同理可发送消息给client1。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 06:34:58