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

使用Riverpod Stream Provider集成SignalR时数据流异常求助

Riverpod Stream Provider 集成SignalR时StreamController无响应、Provider持续加载的问题解决

问题根源分析

你的代码存在几个关键问题,导致StreamController无法正常推送数据、Provider一直处于加载状态:

  1. 重复调用SignalR方法阻塞流程:同时用send和invoke调用了同一个JoinLiveFeed方法,invoke是同步等待返回的调用,会阻塞后续StreamController的初始化和监听逻辑,导致流根本没机会开始接收数据。
  2. 资源未绑定生命周期:没有利用Riverpod的ref.onDispose清理SignalR连接和StreamController,不仅可能导致内存泄漏,也会让流在Provider销毁后无法正确终止。
  3. 类型不匹配:Provider声明返回Stream<Object?>,但StreamController是Stream<List<Object>>,类型不一致可能导致数据无法正确yield。

修正后的代码

@riverpod
Stream<List<Object>> signalRStream(Ref ref) async* {
  Logger.root.level = Level.ALL;
  Logger.root.onRecord.listen((LogRecord rec) {
    print('${rec.level.name}: ${rec.time}: ${rec.message}');
  });

  // 获取token并做空安全判断
  final token = await ref.read(settingsNotifierProvider.notifier).getToken();
  if (token == null) {
    yield [];
    return;
  }

  const String url = 'xxxxxxxxxxxxxx.xxx';
  final hubProtLogger = Logger("SignalR - hub");
  final hubConnection = HubConnectionBuilder()
      .withUrl(url,
          options: HttpConnectionOptions(
              skipNegotiation: true,
              logger: hubProtLogger,
              transport: HttpTransportType.WebSockets,
              requestTimeout: 50000,
              accessTokenFactory: () async => token))
      .build();

  // 绑定资源清理逻辑,Provider销毁时关闭连接和控制器
  ref.onDispose(() async {
    await hubConnection.stop();
    await _controller.close();
  });

  // 启动SignalR连接
  await hubConnection.start();
  print("SignalR connection started");

  // 仅调用一次JoinLiveFeed,根据业务需求选invoke(需等待返回)或send(无需等待)
  await hubConnection.invoke("JoinLiveFeed", args: <Object>[
    [37, 38]
  ]);
  print("Joined live feed successfully");

  final StreamController<List<Object>> _controller = StreamController();

  // 监听SignalR推送的消息,确保控制器未关闭时再添加数据
  hubConnection.on("SendSymbolsListAsync", (parameters) {
    if (!_controller.isClosed) {
      _controller.add(parameters ?? []);
    }
  });

  // 转发流数据到Provider
  await for (final value in _controller.stream) {
    yield value;
  }
}

关键修正点说明

  • 移除重复调用:删掉了冗余的hubConnection.send调用,避免阻塞后续流初始化流程。
  • 添加资源清理:通过ref.onDispose在Provider销毁时自动关闭SignalR连接和StreamController,防止内存泄漏。
  • 统一类型匹配:将Provider返回类型改为Stream<List<Object>>,和StreamController的类型保持一致,确保数据能正确推送。
  • 强化空安全:增加token空判断、控制器状态判断,避免运行时异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 01:20:58