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
相关产品推荐
相关产品推荐

