如何用RXJava 2实现Socket连接建立前的请求排队执行?
如何用RxJava2实现Socket连接就绪后批量处理排队请求
这个场景其实很适合用RxJava的Subject来做请求队列,结合连接状态的流来触发执行,下面我一步步给你拆解实现思路和代码:
核心思路
- 用一个状态流监听Socket的连接状态,仅当状态变为
已连接时才开始处理请求 - 用一个请求队列Subject收集用户发起的所有请求,直到连接就绪
- 连接就绪后,把队列里的请求逐一(或并行)通过Socket发送,并将每个请求的结果单独返回给调用者
- 额外处理连接失败/重连的场景(实际业务中非常实用)
具体实现步骤
1. 定义必要的模型和工具类
先把基础的状态、请求包装、结果类定义好,让逻辑更清晰:
// Socket连接状态枚举 enum ConnectionState { DISCONNECTED, CONNECTING, CONNECTED } // 包装用户请求的类,包含请求数据和结果回调发射器 class SocketRequest { private final String requestData; private final SingleEmitter<Response> emitter; public SocketRequest(String requestData, SingleEmitter<Response> emitter) { this.requestData = requestData; this.emitter = emitter; } // 拿到Socket实例后执行请求的核心逻辑 public void execute(Socket socket) { try { // 实际Socket发送逻辑:写入请求数据 OutputStream outputStream = socket.getOutputStream(); outputStream.write(requestData.getBytes()); outputStream.flush(); // 读取并解析响应 InputStream inputStream = socket.getInputStream(); Response response = parseResponse(inputStream); emitter.onSuccess(response); } catch (IOException e) { emitter.onError(e); } } // 简化的响应解析逻辑,实际业务中替换为真实解析 private Response parseResponse(InputStream inputStream) throws IOException { return new Response("Response for: " + requestData); } } // 请求响应结果类 class Response { private final String content; public Response(String content) { this.content = content; } public String getContent() { return content; } }
2. 构建Socket请求管理类
这个类是核心,负责连接管理、请求队列维护和请求执行:
import io.reactivex.*; import io.reactivex.disposables.Disposable; import io.reactivex.schedulers.Schedulers; import io.reactivex.subjects.BehaviorSubject; import java.net.Socket; import java.util.concurrent.ConcurrentLinkedQueue; public class SocketRequestManager { // 连接状态流:BehaviorSubject保存最新状态,新订阅者能立即拿到当前状态 private final BehaviorSubject<ConnectionState> connectionStateSubject = BehaviorSubject.createDefault(ConnectionState.DISCONNECTED); // 线程安全的请求队列,多线程下也能安全添加/取出请求 private final ConcurrentLinkedQueue<SocketRequest> requestQueue = new ConcurrentLinkedQueue<>(); // 缓存Socket连接实例,避免重复发起连接 private Single<Socket> socketSingle; // 连接的Disposable,用于断开时清理资源 private Disposable connectionDisposable; // 对外暴露的请求入口:用户调用此方法发起请求,返回结果的Single public Single<Response> sendRequest(String requestData) { return Single.create(emitter -> { ConnectionState currentState = connectionStateSubject.getValue(); if (currentState == ConnectionState.CONNECTED) { // 已连接,直接执行请求 executeRequest(new SocketRequest(requestData, emitter)); } else { // 未连接,加入队列 requestQueue.add(new SocketRequest(requestData, emitter)); // 如果还没在连接,触发连接流程 if (currentState == ConnectionState.DISCONNECTED) { startSocketConnection(); } } }).subscribeOn(Schedulers.io()); } // 启动Socket连接的逻辑 private void startSocketConnection() { connectionStateSubject.onNext(ConnectionState.CONNECTING); // 创建连接Single,用cache()避免重复订阅导致多次连接 socketSingle = Single.fromCallable(() -> { // 实际业务中替换为你的Socket连接地址和端口 Socket socket = new Socket("localhost", 8080); socket.setKeepAlive(true); return socket; }).subscribeOn(Schedulers.io()) .doOnSuccess(socket -> { connectionStateSubject.onNext(ConnectionState.CONNECTED); // 连接成功后,批量处理队列里的所有请求 processQueuedRequests(socket); }) .doOnError(error -> { connectionStateSubject.onNext(ConnectionState.DISCONNECTED); // 连接失败,通知所有排队的请求出错 notifyQueueRequestsError(error); }) .cache(); connectionDisposable = socketSingle.subscribe(); } // 处理队列中所有等待的请求 private void processQueuedRequests(Socket socket) { SocketRequest request; while ((request = requestQueue.poll()) != null) { executeRequest(request, socket); } } // 执行单个请求的重载方法 private void executeRequest(SocketRequest request) { socketSingle.subscribe(socket -> executeRequest(request, socket), request.emitter::onError); } private void executeRequest(SocketRequest request, Socket socket) { // 在IO线程执行Socket操作,避免阻塞主线程 Completable.fromAction(() -> request.execute(socket)) .subscribeOn(Schedulers.io()) .subscribe(() -> {}, request.emitter::onError); } // 连接失败/断开时,通知所有排队请求 private void notifyQueueRequestsError(Throwable error) { SocketRequest request; while ((request = requestQueue.poll()) != null) { request.emitter.onError(error); } } // 对外暴露的断开连接方法,清理资源和队列 public void disconnect() { if (connectionDisposable != null && !connectionDisposable.isDisposed()) { connectionDisposable.dispose(); } connectionStateSubject.onNext(ConnectionState.DISCONNECTED); notifyQueueRequestsError(new IllegalStateException("Socket disconnected manually")); } }
3. 实际使用示例
在你的Activity、ViewModel或业务类中这样调用:
SocketRequestManager requestManager = new SocketRequestManager(); // 用户发起第一个请求,此时Socket未连接,自动加入队列 requestManager.sendRequest("User Request 1") .subscribe(response -> { // 处理第一个请求的响应 Log.d("SocketDemo", "Response 1: " + response.getContent()); }, error -> { // 处理请求失败 Log.e("SocketDemo", "Request 1 failed", error); }); // 用户发起第二个请求,同样加入队列 requestManager.sendRequest("User Request 2") .subscribe(response -> { Log.d("SocketDemo", "Response 2: " + response.getContent()); }, error -> { Log.e("SocketDemo", "Request 2 failed", error); }); // 当Socket连接成功后,两个请求会被依次执行,各自返回结果
关键细节说明
- BehaviorSubject:确保新发起的请求能立即获取当前连接状态,避免重复判断
- ConcurrentLinkedQueue:多线程安全的队列,适配用户可能在不同线程发起请求的场景
- Single.cache():缓存Socket连接实例,所有请求共享同一个连接,避免重复创建连接
- SingleEmitter:每个请求绑定独立的发射器,保证调用者能拿到自己请求的专属结果
- 连接失败/断开时的队列清理:避免请求无限等待,及时通知调用者错误信息
内容的提问来源于stack exchange,提问作者Avi Barel
相关产品推荐
相关产品推荐

