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

Flutter中如何在Bloc异步函数运行时重启它?代码问题排查

问题描述

我有一个使用rxdart控制器的简单Bloc(ReaderBloc),仅负责流式输出字符串,内置重启读取的逻辑,代码如下:

class ReaderBloc {
  bool started = false;
  final _publishStream = PublishSubject<String>();
  Stream<String> get publishStream => _publishStream.stream;

  startReading() async {
    List<String> lines = ['a','b','c','d','e','f','g','h'];

    for (String l in lines) {
      if (started) {
        started = false;
        _publishStream.done;
        break;
      }
      _publishStream.add(l);
      await Future.delayed(const Duration(milliseconds: 3000));
    }
  }

  dispose() {
    _publishStream.close();
  }
}

对应的视图组件Reader用于接收并显示Bloc输出的字符串,代码如下:

class Reader extends StatefulWidget {
  const Reader({super.key});
  @override
  State<Reader> createState() => _ReaderState();
}

class _ReaderState extends State<Reader> {
  late ReaderBloc readerBloc;

  @override
  void dispose() {
    super.dispose();
  }

  @override
  Widget build(BuildContext context) {
    readerBloc = Provider.of<ReaderBloc>(context);

    return DefaultTabController(
      length: 1,
      child: Scaffold(
        appBar: AppBar(
          actions: <Widget>[
            IconButton(
              icon: const Icon(
                Icons.refresh,
                color: Colors.white,
              ),
              onPressed: () {
                // seems this part doesn't work
                readerBloc.started = true;
                readerBloc.startReading();
              },
            )
          ],
        ),
        body: TabBarView(
          children: [
            Scaffold(body: tabBody(readerBloc.publishStream)),
          ],
        ),
      ),
    );
  }

  tabBody(Stream<String> stream) {
    return StreamBuilder<String>(
      stream: stream,
      builder: (context, snapshot) {
        if (!snapshot.hasData) {
          return const Center(child: CircularProgressIndicator());
        }
        return Text(
          snapshot.data.toString(),
        );
      },
    );
  }
}

进入Reader路由前,我执行了以下操作:

Provider.of<ReaderBloc>(context, listen: false).startReading();
Navigator.of(context).push(MaterialPageRoute(builder: (context) => const Reader()));

问题:Reader视图中的刷新按钮本应从头开始重新读取,但未达到预期效果——既不停止当前读取,也不启动新的读取,且无任何报错。请问我遗漏了什么?


问题分析与修复

核心问题点

  1. started变量逻辑错误与线程冲突:

    • 初始值false,循环内判断if (started)才停止任务,和“标记重启就终止当前任务”的需求逻辑完全相反。
    • 直接修改布尔变量并调用异步方法,会导致多个startReading()任务同时运行,无同步机制导致状态混乱。
  2. _publishStream.done是无效操作:

    • 该代码只是获取一个Future对象,没有执行任何关闭或重置流的操作,无法终止当前数据发送。
  3. Bloc初始化位置不符合生命周期规范:

    • 在build方法中重复赋值readerBloc,虽不影响功能,但不符合Flutter StatefulWidget的最佳实践。

修复后的代码

1. 修正ReaderBloc逻辑

class ReaderBloc {
  // 用CancelableOperation管理异步任务,实现安全取消
  CancelableOperation? _currentReadingTask;
  final _publishStream = PublishSubject<String>();
  Stream<String> get publishStream => _publishStream.stream;

  startReading() async {
    // 先取消正在运行的旧任务
    _currentReadingTask?.cancel();
    
    List<String> lines = ['a','b','c','d','e','f','g','h'];
    // 创建新的可取消任务
    _currentReadingTask = CancelableOperation.fromFuture(
      Future.forEach(lines, (String l) async {
        _publishStream.add(l);
        await Future.delayed(const Duration(milliseconds: 3000));
      }),
    );

    // 任务完成或取消后清空引用
    await _currentReadingTask?.valueOrCancellation();
    _currentReadingTask = null;
  }

  dispose() {
    _currentReadingTask?.cancel();
    _publishStream.close();
  }
}

2. 修正Reader组件逻辑

class Reader extends StatefulWidget {
  const Reader({super.key});
  @override
  State<Reader> createState() => _ReaderState();
}

class _ReaderState extends State<Reader> {
  late ReaderBloc readerBloc;

  @override
  void initState() {
    super.initState();
    // 在initState中初始化Bloc,避免build重复赋值
    readerBloc = Provider.of<ReaderBloc>(context, listen: false);
  }

  @override
  void dispose() {
    super.dispose();
  }

  @override
  Widget build(BuildContext context) {
    return DefaultTabController(
      length: 1,
      child: Scaffold(
        appBar: AppBar(
          actions: <Widget>[
            IconButton(
              icon: const Icon(
                Icons.refresh,
                color: Colors.white,
              ),
              onPressed: () {
                // 直接调用startReading,内部已处理旧任务取消
                readerBloc.startReading();
              },
            )
          ],
        ),
        body: TabBarView(
          children: [
            Scaffold(body: tabBody(readerBloc.publishStream)),
          ],
        ),
      ),
    );
  }

  Widget tabBody(Stream<String> stream) {
    return StreamBuilder<String>(
      stream: stream,
      builder: (context, snapshot) {
        if (!snapshot.hasData) {
          return const Center(child: CircularProgressIndicator());
        }
        return Text(
          snapshot.data.toString(),
        );
      },
    );
  }
}

关键修复说明

  • 用CancelableOperation替代布尔变量,安全取消正在运行的异步任务,避免多任务冲突。
  • 每次调用startReading()先终止旧任务再启动新任务,确保重启逻辑可靠。
  • 将Bloc初始化移到initState,符合Flutter组件生命周期规范。
  • 移除无效的_publishStream.done操作,改用任务取消终止数据发送。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 05:17:08