如何使用CompletableFuture实现Java客户端等待WebRTC应用的RPC服务器响应
Great question! Ditching Thread.sleep() and manual latches for CompletableFuture is exactly the right move here—it’ll keep your client responsive while still letting you wait for the RPC response in a clean, non-blocking way. Here’s a step-by-step solution tailored to your code:
Core Idea
We’ll create a thread-safe map to track pending RPC requests by their UUID. When you send a request, you’ll create a CompletableFuture, store it in the map with the request’s ID, then return the future immediately. When your WebSocket receives a response from the server, you’ll look up the matching future in the map and complete it with the server’s data (or an error if something went wrong).
Step 1: Add a Pending Requests Map
First, add a thread-safe map as a member variable in your client class to track in-flight requests. ConcurrentHashMap is perfect here because it handles concurrent access safely:
private final ConcurrentHashMap<String, CompletableFuture<SessionDescription>> pendingRequests = new ConcurrentHashMap<>(); private final ObjectMapper objectMapper = new ObjectMapper(); // Assuming this is already a member
Step 2: Update the join Method
Modify your join method to create a CompletableFuture, register it in the map, send the RPC request, and return the future. This way, the method doesn’t block—callers can use thenAccept(), join(), or other CompletableFuture methods to handle the response when it arrives:
@Override public CompletableFuture<SessionDescription> join(String sid, String uid, SessionDescription offer) { String uuid = UUID.randomUUID().toString(); CompletableFuture<SessionDescription> future = new CompletableFuture<>(); // Register the future with the request ID pendingRequests.put(uuid, future); JsonRpcRequestMessage rpcMsg = new JsonRpcRequestMessage(); rpcMsg.setJsonrpc("2.0"); rpcMsg.setId(uuid); rpcMsg.setMethod("join"); Map<String, Object> params = new HashMap<>(); params.put("offer", offer); params.put("sid", sid); params.put("uid", uid); rpcMsg.setParams(params); try { String rpcText = objectMapper.writeValueAsString(rpcMsg); webSocket.send(rpcText); } catch (JsonProcessingException e) { // If sending fails, complete the future with an exception and clean up future.completeExceptionally(e); pendingRequests.remove(uuid); e.printStackTrace(); } // Return the future immediately—callers can handle the response later return future; }
Step 3: Handle WebSocket Responses
Next, you need to hook into your WebSocket’s message-received callback. When the server sends back a JSON-RPC response, parse it, find the matching future in the map, and complete it. Here’s how that might look:
// This method should be called whenever your WebSocket receives a message private void handleWebSocketMessage(String message) { try { // Parse the JSON-RPC response (you'll need a response DTO) JsonRpcResponseMessage response = objectMapper.readValue(message, JsonRpcResponseMessage.class); String requestId = response.getId(); // Look up the pending future CompletableFuture<SessionDescription> future = pendingRequests.remove(requestId); if (future != null) { if (response.getError() != null) { // If the server returned an error, complete the future with an exception future.completeExceptionally(new RuntimeException("RPC error: " + response.getError().getMessage())); } else { // Parse the successful response into SessionDescription and complete the future SessionDescription result = objectMapper.convertValue(response.getResult(), SessionDescription.class); future.complete(result); } } } catch (JsonProcessingException e) { e.printStackTrace(); // Optionally handle invalid messages here } }
You’ll need a simple DTO for the JSON-RPC response, like this:
public class JsonRpcResponseMessage { private String jsonrpc; private String id; private Object result; private JsonRpcError error; // Getters and setters public static class JsonRpcError { private int code; private String message; // Getters and setters } }
Step 4: Add Timeout Protection (Optional but Recommended)
To prevent futures from hanging indefinitely if the server never responds, you can wrap the returned future with a timeout. Modify the join method’s return line to:
return future.orTimeout(10, TimeUnit.SECONDS) // Adjust timeout as needed .whenComplete((res, err) -> pendingRequests.remove(uuid)); // Clean up even on timeout
This will complete the future with a TimeoutException if no response arrives within 10 seconds, and ensure we clean up the pending request entry.
Why This Works Better
- Non-blocking: Your client thread doesn’t wait or sleep—other tasks can run while waiting for the RPC response.
- Responsive: The future completes as soon as the server sends a response, no guesswork about sleep durations.
- Clean error handling: You can propagate RPC errors and network failures directly through the future’s exception chain.
内容的提问来源于stack exchange,提问作者user9383769

