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

Flutter中接收WebSocket响应能否不用Stream改用Future实现

Flutter WebSocket 使用Future等待服务端响应的实现方案

首先明确:你要的这种调用方式是可以实现的,不过本质是对WebSocket原生Stream API做了上层封装,完全可以做到上层调用时只通过await拿单次请求的响应。

实现原理

WebSocket本身是全双工长连接协议,原生API暴露Stream是为了适配服务端可以主动推送多条消息的特性,而Future只能返回单次异步结果,所以我们可以基于Completer实现「单次请求-单次响应」的Future封装,核心逻辑是给每个请求生成唯一标识,响应时根据标识匹配对应的Future返回结果。

代码实现示例

1. 封装WebSocket工具类

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

class WebSocketClient {
  final WebSocketChannel channel;
  // 存储请求ID与对应Completer的映射
  final Map<String, Completer> _pendingRequests = {};
  int _reqIdSeq = 0;

  WebSocketClient(this.channel) {
    _initMessageListener();
  }

  void _initMessageListener() {
    channel.stream.listen(
      (rawMessage) {
        // 假设服务端和客户端约定消息为JSON格式,携带reqId匹配请求
        final message = jsonDecode(rawMessage);
        final reqId = message['reqId'] as String?;
        if (reqId != null && _pendingRequests.containsKey(reqId)) {
          final completer = _pendingRequests.remove(reqId)!;
          if (!completer.isCompleted) {
            completer.complete(message['data']);
          }
        }
        // 服务端主动推送的无reqId消息可在此处单独处理
      },
      onError: (error) {
        // 连接出错时终止所有等待中的请求
        for (final completer in _pendingRequests.values) {
          if (!completer.isCompleted) completer.completeError(error);
        }
        _pendingRequests.clear();
      },
      onDone: () {
        // 连接断开时终止所有等待中的请求
        for (final completer in _pendingRequests.values) {
          if (!completer.isCompleted) completer.completeError(Exception('WebSocket连接已断开'));
        }
        _pendingRequests.clear();
      },
    );
  }

  // 对外暴露的带Future返回的请求方法
  Future<T> sendRequest<T>(dynamic requestData) async {
    final reqId = (_reqIdSeq++).toString();
    final completer = Completer<T>();
    _pendingRequests[reqId] = completer;

    // 发送携带reqId的请求
    channel.sink.add(jsonEncode({
      'reqId': reqId,
      'data': requestData,
    }));

    // 增加超时逻辑,避免请求无限等待
    return completer.future.timeout(
      const Duration(seconds: 15),
      onTimeout: () {
        _pendingRequests.remove(reqId);
        throw TimeoutException('请求超时');
      },
    );
  }

  // 销毁方法,页面退出时调用
  void dispose() {
    channel.sink.close();
  }
}

2. 调用方式(符合你期望的写法)

// 初始化WebSocket连接
final wsChannel = WebSocketChannel.connect(Uri.parse('ws://你的服务端接口地址'));
final wsClient = WebSocketClient(wsChannel);

try {
  // 直接await拿到响应
  final response = await wsClient.sendRequest({'action': 'getUserInfo', 'userId': 1001});
  // 处理拿到的响应数据
  print('服务端返回结果:$response');
} catch (e) {
  // 处理超时、连接错误等异常
  print('请求出错:$e');
}

注意事项

  • 该方案仅适用于客户端发请求、服务端返回对应响应的一对一匹配场景,如果是服务端主动推送的通知类消息,还是需要单独用回调或者Stream处理
  • 服务端必须配合在响应中返回和请求一致的reqId,否则会出现请求匹配错误
  • 记得在页面/应用退出时调用dispose方法关闭连接,避免内存泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 12:45:03