【问题标题】:Implementing web socket channel on flutter and getting Bad state : stream has already been listened when calling the event again在 flutter 上实现 web 套接字通道并获得错误状态:再次调用事件时已经侦听了流
【发布时间】:2023-01-31 02:37:40
【问题描述】:

当我在应用程序从非活动状态恢复后调用 bloc 时,我遇到了上述问题,以便我从 websocket 获得新的数据流。

我已经与大多数说明共享了我的代码块。

问题出在 websocket 流订阅管理的某个地方,我已经尝试了很多,但在某些时候我被阻止了,无法继续进行

///Important///
///---packages needed
///web_socket_channel
///flutter_bloc
///---for multi bloc provider
///BlocProvider(create: (context) => SampleStreamBloc()),
///---reproduce the issue
///call the bloc with event as below from more than one screen
///context.read<SampleStreamBloc>().add(SampleStreamConnect());

import 'dart:convert';
import 'dart:developer';

import 'package:flutter_bloc/flutter_bloc.dart';
import 'package:web_socket_channel/io.dart';

IOWebSocketChannel channelStocks = IOWebSocketChannel.connect(Uri.parse('wss://ws.eodhistoricaldata.com/ws/forex?api_token=demo'));

class SampleStreamBloc extends Bloc<SampleStreamEvent, SampleStreamState> {
  SampleStreamBloc() : super(SampleStreamInitial()) {
   
    //event for listening from socket
    on<SampleStreamConnect>((event, emit) async {
      emit(SampleStreamProgress());
      await emit.forEach(channelStocks.stream, onData: ((data) {
        Map<String, dynamic> message = jsonDecode(data);
        log(message.toString(), name: "Stream Response");
        if (message['message'] == "Authorized") {
          var data = jsonEncode({"action": "subscribe", "symbols": "EURUSD"});
          channelStocks.sink.add(data);
        }
        return SampleStreamSuccess();
      }), onError: (e, stackTrace) {
        log(e.toString(), error: e, stackTrace: stackTrace);
        return SampleStreamSuccess();
      });
    });
  }
}

//states
abstract class SampleStreamState {}

//events
abstract class SampleStreamEvent {}

//states implementation
class SampleStreamInitial extends SampleStreamState {}

class SampleStreamProgress extends SampleStreamState {}

class SampleStreamSuccess extends SampleStreamState {}

//events implementation
class SampleStreamConnect extends SampleStreamEvent {}


也分享了错误。

【问题讨论】:

  • 这意味着您正在代码中的其他地方调用 channelStocks.stream.listen 方法(直接或间接)
  • 不,但我再次调用 bloc 的 SampleStreamConnect 事件,因为当数据流由于任何原因停止到一定持续时间时,我需要重新连接到同一个套接字并且数据需要在同一个 bloc @pskink 中发出
  • 如果有任何其他方法可以实现我在上面评论中描述的内容,请提出建议!

标签: flutter bloc


【解决方案1】:

我建议你使用这个多平台 websocket 包https://pub.dev/packages/websocket_universal,在那里你甚至可以跟踪所有发生的 WS 事件(如果你需要的话,甚至可以内置 ping 测量):

import 'package:websocket_universal/websocket_universal.dart';

/// Example works with Postman Echo server
void main() async {
  /// Postman echo ws server (you can use your own server URI)
  /// 'wss://ws.postman-echo.com/raw'
  /// For local server it could look like 'ws://127.0.0.1:42627/websocket'
  const websocketConnectionUri = 'wss://ws.postman-echo.com/raw';
  const textMessageToServer = 'Hello server!';
  const connectionOptions = SocketConnectionOptions(
    pingIntervalMs: 3000, // send Ping message every 3000 ms
    timeoutConnectionMs: 4000, // connection fail timeout after 4000 ms
    /// see ping/pong messages in [logEventStream] stream
    skipPingMessages: false,

    /// Set this attribute to `true` if do not need any ping/pong
    /// messages and ping measurement. Default is `false`
    pingRestrictionForce: false,
  );

  /// Example with simple text messages exchanges with server
  /// (not recommended for applications)
  /// [<String, String>] generic types mean that we receive [String] messages
  /// after deserialization and send [String] messages to server.
  final IMessageProcessor<String, String> textSocketProcessor =
      SocketSimpleTextProcessor();
  final textSocketHandler = IWebSocketHandler<String, String>.createClient(
    websocketConnectionUri, // Postman echo ws server
    textSocketProcessor,
    connectionOptions: connectionOptions,
  );

  // Listening to webSocket status changes
  textSocketHandler.socketHandlerStateStream.listen((stateEvent) {
    // ignore: avoid_print
    print('> status changed to ${stateEvent.status}');
  });

  // Listening to server responses:
  textSocketHandler.incomingMessagesStream.listen((inMsg) {
    // ignore: avoid_print
    print('> webSocket  got text message from server: "$inMsg" '
        '[ping: ${textSocketHandler.pingDelayMs}]');
  });

  // Listening to debug events inside webSocket
  textSocketHandler.logEventStream.listen((debugEvent) {
    // ignore: avoid_print
    print('> debug event: ${debugEvent.socketLogEventType}'
        ' [ping=${debugEvent.pingMs} ms]. Debug message=${debugEvent.message}');
  });

  // Listening to outgoing messages:
  textSocketHandler.outgoingMessagesStream.listen((inMsg) {
    // ignore: avoid_print
    print('> webSocket sent text message to   server: "$inMsg" '
        '[ping: ${textSocketHandler.pingDelayMs}]');
  });

  // Connecting to server:
  final isTextSocketConnected = await textSocketHandler.connect();
  if (!isTextSocketConnected) {
    // ignore: avoid_print
    print('Connection to [$websocketConnectionUri] failed for some reason!');
    return;
  }

  textSocketHandler.sendMessage(textMessageToServer);

  await Future<void>.delayed(const Duration(seconds: 30));
  // Disconnecting from server:
  await textSocketHandler.disconnect('manual disconnect');
  // Disposing webSocket:
  textSocketHandler.close();
}

【讨论】:

    猜你喜欢
    • 2023-02-16
    • 2020-02-18
    • 2011-12-10
    • 2018-09-01
    • 1970-01-01
    • 2022-01-22
    • 2020-03-08
    • 2021-06-10
    • 2021-12-26
    相关资源
    最近更新 更多