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

Dart中WebSocket请求如何等待对应响应?

解决WebSocket请求-响应同步等待问题

问题背景

使用Dart的IOWebSocketChannel建立WebSocket连接,初始化时需要循环发送多个请求并处理对应响应。当前采用将响应存入列表后轮询的方式,出现超时后才批量收到响应的问题,无法改用HTTP请求,且为非Flutter项目,需寻找可行的等待响应方案。

当前现象

实际打印输出

send data and wait for response
send data and wait for response
send data and wait for response
send data and wait for response
send data and wait for response
new webSocket Message: ...     // event message simplified with "..."
new webSocket Message: ...
new webSocket Message: ...
new webSocket Message: ...
new webSocket Message: ...

期望输出

send data and wait for response
new webSocket Message: ...
data found
send data and wait for response
new webSocket Message: ...
data found
send data and wait for response
new webSocket Message: ...
data found
send data and wait for response
new webSocket Message: ...
data found

当前简化代码

class Connection {

    late IOWebSocketChannel webSocketChannel;
    var rxData = [];

    connect() {    
        webSocketChannel = IOWebSocketChannel.connect(  
            Uri.parse("ipAdress"),  
            pingInterval: Duration(seconds: 30),  
        );  
      
        webSocketChannel.stream.listen((event) {  
            print("new webSocket Message: $event");  
            rxData.add(event);  
        });   
    }

    send(opcode, parameters) {  
        Map<String, dynamic> dataSendMap = {  
            "opcode": opcode,  
            "parameters": parameters,  
        };  
        String jsonData = json.encode(dataSendMap);  
        webSocketChannel.sink.add(jsonData);  
    }

    sendAndWaitForAnswer(dataOpt, dataParam) {
        print("send data and wait for response");
        send(dataOpt, dataParam);
        timeout = 5;
        int startTime = DateTime.now().millisecondsSinceEpoch ~/ Duration.millisecondsPerSecond;
        while((DateTime.now().millisecondsSinceEpoch ~/ Duration.millisecondsPerSecond) - startTime < timeout) {
            for (var data in rxData) {  
                if (data['opcode'] == dataOpt) {  
                    print("data found");
                    answer = data;  
                    rxData.remove(data);  
                }
            }
        }
        return answer;
    }
}

解决方案

核心问题是同步死循环阻塞了Dart事件循环,导致WebSocket的消息回调无法及时处理收到的响应,所有消息只能在循环结束后批量触发。需改用异步等待机制,避免阻塞事件循环。

优化方案:用Completer+请求唯一标识跟踪响应

为每个请求生成唯一ID,通过Completer异步等待对应响应,同时维护等待请求映射表,收到响应后立即匹配并完成等待。

修改后的代码示例:

import 'dart:async';
import 'dart:convert';
import 'package:web_socket_channel/io.dart';

class Connection {
  late IOWebSocketChannel webSocketChannel;
  // 存储等待中的请求:key为请求唯一ID,value为对应的Completer
  final Map<String, Completer<dynamic>> _pendingRequests = {};
  // 生成唯一ID的计数器(也可使用UUID库替代)
  int _requestIdCounter = 0;

  void connect() {
    webSocketChannel = IOWebSocketChannel.connect(
      Uri.parse("ipAdress"),
      pingInterval: const Duration(seconds: 30),
    );

    webSocketChannel.stream.listen((event) {
      print("new webSocket Message: $event");
      final response = json.decode(event);
      // 从响应中取出请求唯一ID(需要服务端配合返回该字段)
      final requestId = response['requestId'];
      
      if (requestId != null && _pendingRequests.containsKey(requestId)) {
        // 匹配到对应请求,完成等待并清理记录
        _pendingRequests[requestId]!.complete(response);
        _pendingRequests.remove(requestId);
        print("data found");
      } else {
        // 处理非请求触发的推送消息
        print("received unsolicited message: $response");
      }
    }, onError: (error) {
      // 连接出错时,批量完成所有等待中的请求
      for (var completer in _pendingRequests.values) {
        completer.completeError(error);
      }
      _pendingRequests.clear();
    });
  }

  void _send(String requestId, String opcode, dynamic parameters) {
    final dataSendMap = {
      "requestId": requestId,
      "opcode": opcode,
      "parameters": parameters,
    };
    final jsonData = json.encode(dataSendMap);
    webSocketChannel.sink.add(jsonData);
  }

  Future<dynamic> sendAndWaitForAnswer(
    String opcode, 
    dynamic parameters, 
    {Duration timeout = const Duration(seconds: 5)}
  ) async {
    print("send data and wait for response");
    // 生成唯一请求ID
    final requestId = 'req_${_requestIdCounter++}';
    final completer = Completer<dynamic>();
    _pendingRequests[requestId] = completer;

    // 发送请求
    _send(requestId, opcode, parameters);

    // 设置超时逻辑
    final timer = Timer(timeout, () {
      if (!completer.isCompleted) {
        completer.completeError(TimeoutException('Request timed out', timeout));
        _pendingRequests.remove(requestId);
        print("request timed out");
      }
    });

    try {
      return await completer.future;
    } finally {
      timer.cancel();
    }
  }
}

关键说明

  1. 不阻塞事件循环:用Future+Completer实现异步等待,WebSocket的消息回调能实时处理收到的响应,不会被循环阻塞。
  2. 请求-响应精准匹配:通过唯一requestId绑定请求和响应,避免同opcode的多请求匹配错误。
  3. 超时与错误处理:内置超时逻辑,连接出错时自动清理所有等待请求,防止内存泄漏。
  4. 服务端配合:需要服务端在响应中返回请求时的requestId字段,否则无法完成匹配。

循环发送示例

void main() async {
  final conn = Connection();
  conn.connect();
  // 等待连接建立(可根据实际情况调整延迟或监听连接状态)
  await Future.delayed(const Duration(seconds: 1));

  // 循环发送请求并逐个等待响应
  for (var i = 0; i < 4; i++) {
    try {
      final response = await conn.sendAndWaitForAnswer('some_opcode', {'param': i});
      print('processed response: $response');
    } catch (e) {
      print('request failed: $e');
    }
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 15:55:06