如何基于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("无效指令,请重新输入"); } } } }
五、测试流程
- 启动
ChatServer; - 启动第一个客户端,输入用户ID(如
client1)完成登录; - 启动第二个客户端,输入用户ID(如
client2)完成登录; - 在
client1控制台输入指令1,输入目标用户IDclient2发起房间创建请求; - 两个客户端会收到服务器推送的房间配置信息,自动记录房间ID;
- 在
client1输入指令2,输入消息内容,client2会实时收到该消息; client2同理可发送消息给client1。
内容的提问来源于stack exchange,提问作者Svetlana
相关产品推荐
相关产品推荐

