如何在BLOC中取消订阅流的异步操作?
问题描述
我有一个BLOC,其中包含异步函数,会根据用例事件流生成事件:
on<MyBlocEvent>((event, emit) async { stream = useCaseResponse.stream; await stream.forEach((MyModel event) { emit(state.copyWith(property: event.property)); }); });
当用户在前端点击“取消”时,会触发另一个事件MyBlocCancelEvent,需要取消所有流相关操作并emit初始状态:
on<MyBlocCancelEvent>((event, emit) async { // 需要取消所有流操作 // 同时emit初始状态,方便用户再次点击按钮重新开始 });
我想到两种实现方式,但都不符合规范:
- 将流保存为BLOC实例属性,在另一事件处理方法中获取
- 将流保存到状态中,通过对应事件取消
请问如何通过另一个事件取消forEach操作?这类需求的标准实现方式是什么?
补充代码
以下是BLOC中的流处理代码:
Stream<PositionDeterminationDto> stream = response.positionStreamDeterminationStream; await stream .forEach((PositionDeterminationDto positionDeterminationDto) { emit(state.copyWith( currentGpsLocationModel: CurrentGpsLocationModel.fromPositionDto( positionDeterminationDto.positionDto), status: DetermineLocationStatus.determineLocationProcessing, )); }).whenComplete(() { emit(state.copyWith( status: DetermineLocationStatus.determineLocationFinished, progress: 0)); }); });
若将Stream.forEach改为Stream.listen,BLOC会报错:
emit was called after an event handler completed normally. This is usually due to an unawaited future in an event handler. Please make sure to await all asynchronous operations with event handlers and use emit.isDone after asynchronous operations before calling emit() to ensure the event handler has not completed.
标准实现方案
在Bloc中处理这类可取消的流操作,标准做法是在Bloc类中维护一个StreamSubscription类型的实例变量,用来保存流的订阅对象,这样就能在取消事件中调用cancel()方法终止流的监听,同时重置状态。
具体步骤:
- 在Bloc类中声明一个
StreamSubscription?类型的实例变量,用来存储流的订阅:
class YourBloc extends Bloc<YourEvent, YourState> { StreamSubscription? _locationStreamSubscription; YourBloc() : super(YourInitialState()) { on<MyBlocEvent>(_handleMyBlocEvent); on<MyBlocCancelEvent>(_handleMyBlocCancelEvent); } }
- 在
MyBlocEvent的处理方法中,初始化订阅并保存到实例变量,同时处理流的事件:
void _handleMyBlocEvent(MyBlocEvent event, Emitter<YourState> emit) async { // 先取消之前可能存在的订阅,避免重复监听 await _locationStreamSubscription?.cancel(); final stream = response.positionStreamDeterminationStream; _locationStreamSubscription = stream.listen((PositionDeterminationDto dto) { emit(state.copyWith( currentGpsLocationModel: CurrentGpsLocationModel.fromPositionDto(dto.positionDto), status: DetermineLocationStatus.determineLocationProcessing, )); }) ..onDone(() { emit(state.copyWith( status: DetermineLocationStatus.determineLocationFinished, progress: 0, )); // 完成后清空订阅 _locationStreamSubscription = null; }); // 等待订阅完成,避免Bloc报错 await _locationStreamSubscription?.asFuture(); }
- 在
MyBlocCancelEvent的处理方法中,取消订阅并重置为初始状态:
void _handleMyBlocCancelEvent(MyBlocCancelEvent event, Emitter<YourState> emit) async { // 取消流订阅 await _locationStreamSubscription?.cancel(); _locationStreamSubscription = null; // 发射初始状态 emit(YourInitialState()); }
关键注意点:
- 每次启动新的流监听前,先取消之前的订阅,防止内存泄漏和重复emit状态
- 使用
asFuture()将StreamSubscription转为Future并await,避免出现emit was called after an event handler completed normally的报错——因为这样能保证事件处理方法在流完成或取消前不会提前结束 - 在流完成或取消后,记得清空订阅变量,避免持有无用引用
内容的提问来源于stack exchange,提问作者harrow
相关产品推荐
相关产品推荐

