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

如何用RXJava 2实现Socket连接建立前的请求排队执行?

如何用RxJava2实现Socket连接就绪后批量处理排队请求

这个场景其实很适合用RxJava的Subject来做请求队列,结合连接状态的流来触发执行,下面我一步步给你拆解实现思路和代码:

核心思路

  1. 用一个状态流监听Socket的连接状态,仅当状态变为已连接时才开始处理请求
  2. 用一个请求队列Subject收集用户发起的所有请求,直到连接就绪
  3. 连接就绪后,把队列里的请求逐一(或并行)通过Socket发送,并将每个请求的结果单独返回给调用者
  4. 额外处理连接失败/重连的场景(实际业务中非常实用)

具体实现步骤

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:26:58