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

Dart中如何管理StreamSubscription多数据 解决Socket重复监听报错

问题原因

Dart 中的 Socket 属于单订阅流,默认同一时间只能注册一个监听器,你代码中连续调用四次 socket.listen() 方法,第二次调用时就会触发你遇到的 Bad state: Stream has already been listened to. 异常。
同时即使没有异常,你的写法逻辑也无法正常运行:多个监听器会乱序接收数据,根本没法保证收到的第一份数据是请求类型、第二份是客户端公钥的约定顺序。

解决方案

要实现按需按顺序等待接收数据,有两种常用实现方式:

方案1:使用 StreamQueue(推荐,逻辑最清晰)

首先在 pubspec.yaml 中引入 Dart 官方的 async 工具包:

dependencies:
  async: ^2.11.0

修改 _handleClient 代码如下:

import 'package:async/async.dart';

void _handleClient(Socket socket) async {
  // 把Socket流包装为支持按顺序读取的StreamQueue
  final socketQueue = StreamQueue<List<int>>(socket);
  late String request, username, password;

  // 按需等待第一条数据:请求类型
  final requestData = await socketQueue.next;
  request = String.fromCharCodes(requestData);

  // 向客户端发送公钥
  socket.write(_publicKey);

  // 按需等待第二条数据:加密的客户端公钥
  final clientPubKeyData = await socketQueue.next;
  _clientPublicKey = _decrypt(String.fromCharCodes(clientPubKeyData));

  // 按需等待第三条数据:加密的用户名
  final usernameData = await socketQueue.next;
  username = _decrypt(String.fromCharCodes(usernameData));

  // 按需等待第四条数据:加密的密码
  final passwordData = await socketQueue.next;
  password = _decrypt(String.fromCharCodes(passwordData));

  var database = Database(username, password);

  switch (request) {
    case 'getItem':
      getItem(socket, database);
      break;
    case 'addItem':
      addItem(socket, database);
      break;
    // 其他业务逻辑
  }
  // 处理完成后关闭资源
  await socketQueue.cancel();
  await socket.close();
}

方案2:转广播流配合first等待(无需额外依赖)

如果不想引入第三方包,可以将单订阅的Socket转为多订阅广播流,每次需要接收数据时等待first属性返回即可:

void _handleClient(Socket socket) async {
  // 转为多订阅广播流
  final broadcastSocket = socket.asBroadcastStream();
  late String request, username, password;

  // 按需等待第一条数据:请求类型
  final requestData = await broadcastSocket.first;
  request = String.fromCharCodes(requestData);

  // 向客户端发送公钥
  socket.write(_publicKey);

  // 按需等待第二条数据:加密的客户端公钥
  final clientPubKeyData = await broadcastSocket.first;
  _clientPublicKey = _decrypt(String.fromCharCodes(clientPubKeyData));

  // 按需等待第三条数据:加密的用户名
  final usernameData = await broadcastSocket.first;
  username = _decrypt(String.fromCharCodes(usernameData));

  // 按需等待第四条数据:加密的密码
  final passwordData = await broadcastSocket.first;
  password = _decrypt(String.fromCharCodes(passwordData));

  var database = Database(username, password);

  switch (request) {
    case 'getItem':
      getItem(socket, database);
      break;
    case 'addItem':
      addItem(socket, database);
      break;
    // 其他业务逻辑
  }
  await socket.close();
}
注意事项
  • 以上实现都依赖客户端严格按照约定顺序发送数据:先发送请求类型,收到服务端公钥后再依次发送加密的客户端公钥、用户名、密码,顺序错误会导致数据解析完全混乱。
  • 生产环境使用建议处理TCP粘包问题:给每个数据包加上固定长度的头部标识数据包长度,接收时先读取头部拿到包长,再按长度读取完整数据包,避免一次收到多份数据或者数据不全的情况。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 00:09:02