Flutter中Bloc+StatelessWidget下如何维持RxDart广播流监听?
交易管理应用Flutter Bloc流订阅问题及解决方案
问题背景
正在学习Flutter和flutter_bloc:8.1.6,开发一款简单的交易管理应用,支持用户添加/编辑交易,要求新增或更新交易后交易列表能实时更新。
底层通过RxDart的BehaviorSubject暴露数据流,相关代码如下:
仓库层代码
class TransactionRepository { final LocalDB _database; final BehaviorSubject<List<Transaction>> _transactionsController = BehaviorSubject.seeded([]); TransactionRepository({required LocalDB database}) : _database = database { _init(); } Future<void> _init() async { try { final transactions = await _getAllTransactions(); _transactionsController.add(transactions); } catch (e) { rethrow; } } Stream<List<Transaction>> streamTransactions() { return _transactionsController.asBroadcastStream(); } /// Save a transaction Future<void> createTransaction(NewTransaction transaction) async { try { await _database .into(_database.transactionTable) .insert(TransactionTableCompanion( //... fields )); final transactions = await _getAllTransactions(); _transactionsController.add(transactions); } catch (e) { rethrow; } } }
Bloc层代码
class TransactionsBloc extends Bloc<TransactionsEvent, TransactionsState> { final TransactionRepository _transactionRepository; TransactionsBloc({required TransactionRepository transactionRepository}) : _transactionRepository = transactionRepository, super(const TransactionsState( transactionStatus: TransactionsStatus.initial)) { on<TransactionsEvent>( (event, emit) => event.map( loadAll: (event) => _loadTransactions(event, emit), error: (event) => _errorLoadingTransactions(event, emit), ), transformer: null, ); } Future<void> _loadTransactions( TransactionsEvent event, Emitter<TransactionsState> emit) async { try { emit(state.copyWith(transactionStatus: TransactionsStatus.fetching)); await emit.forEach( _transactionRepository.streamTransactions(), onData: (transactions) { if (transactions.isEmpty) { return state.copyWith( transactionStatus: TransactionsStatus.initial, transactions: [], ); } return state.copyWith( transactionStatus: TransactionsStatus.fetchedSuccessfully, transactions: transactions, ); }, onError: (e, s) { return state.copyWith( transactionStatus: TransactionsStatus.fetchingFailed, message: 'Error loading transaction(s)', ); }, ); } catch (e) { add(const TransactionsEvent.error('Error loading transaction(s)')); } } }
Bloc注入代码(main.dart)
return CupertinoApp( title: appTitle, home: SafeArea( child: RepositoryProvider( create: (BuildContext context) { final database = context.read<DatabaseCubit>(); return TransactionRepository( database: database.state, ); }, child: BlocProvider( create: (BuildContext context) => TransactionsBloc( transactionRepository: context.read<TransactionRepository>(), )..add(const TransactionsEvent.loadAll()), // stream is subscribed here child: const HomeScreen(title: appTitle), ), ), ), );
页面及列表展示代码
HomeScreen:
@override Widget build(BuildContext context) { return CupertinoPageScaffold( navigationBar: CupertinoNavigationBar( leading: const Icon(CupertinoIcons.line_horizontal_3), trailing: IconButton( onPressed: () { Navigator.push( context, CupertinoPageRoute( builder: (BuildContext context) => const AddNewTransaction(), ), ); }, icon: const Icon(CupertinoIcons.add), ), ), child: const TransactionsList(), ); }
TransactionsList:
BlocBuilder<TransactionsBloc, TransactionsState>( builder: (BuildContext context, TransactionsState state) { switch (state.transactionStatus) { case TransactionsStatus.initial: return TransactionListItem( transaction: Transaction.nullTransaction); case TransactionsStatus.fetching: return const Center( child: CupertinoActivityIndicator(), ); case TransactionsStatus.fetchedSuccessfully: return GroupedTransactionsListView( transactions: state.transactions ?? [Transaction.nullTransaction], ); case TransactionsStatus.fetchingFailed: return Center( child: Text(state.message ?? 'Error loading transaction(s)'), ); default: return const Center( child: Text('Something went wrong! :('), ); } }, ),
当前问题
应用启动时流订阅成功,能正常显示所有交易数据;但用户导航到AddTransaction页面添加交易并保存成功后返回HomeScreen,交易列表未更新。调试发现流监听(_loadTransactions方法)已关闭。
需要解决:
- 在StatelessWidget架构下,如何维持流不关闭?
- 若无法维持,返回交易列表页时重新订阅流的推荐方案是什么?
解决方案
方案一:修正Bloc流订阅逻辑,维持长订阅
问题核心是emit.forEach会在Stream完成时结束Future,导致Bloc的事件处理流程终止,后续流数据无法触发状态更新。改用手动管理StreamSubscription的方式维持长订阅:
class TransactionsBloc extends Bloc<TransactionsEvent, TransactionsState> { final TransactionRepository _transactionRepository; StreamSubscription? _transactionsSubscription; // 新增订阅管理变量 TransactionsBloc({required TransactionRepository transactionRepository}) : _transactionRepository = transactionRepository, super(const TransactionsState(transactionStatus: TransactionsStatus.initial)) { on<TransactionsEvent>( (event, emit) => event.map( loadAll: (event) => _loadTransactions(event, emit), error: (event) => _errorLoadingTransactions(event, emit), ), transformer: null, ); } Future<void> _loadTransactions( TransactionsEvent event, Emitter<TransactionsState> emit) async { try { emit(state.copyWith(transactionStatus: TransactionsStatus.fetching)); // 先取消已有订阅,避免重复订阅导致内存泄漏 await _transactionsSubscription?.cancel(); // 重新订阅流并保存订阅实例 _transactionsSubscription = _transactionRepository.streamTransactions().listen( (transactions) { if (transactions.isEmpty) { emit(state.copyWith( transactionStatus: TransactionsStatus.initial, transactions: [], )); } else { emit(state.copyWith( transactionStatus: TransactionsStatus.fetchedSuccessfully, transactions: transactions, )); } }, onError: (e, s) { emit(state.copyWith( transactionStatus: TransactionsStatus.fetchingFailed, message: 'Error loading transaction(s)', )); }, ); } catch (e) { add(const TransactionsEvent.error('Error loading transaction(s)')); } } // 重写close方法,取消订阅并关闭Repository的Subject,防止内存泄漏 @override Future<void> close() { _transactionsSubscription?.cancel(); _transactionRepository._transactionsController.close(); return super.close(); } }
方案二:返回列表页时重新触发流订阅
若不想维持长订阅,可在从AddTransaction页面返回后,手动触发loadAll事件重新订阅流:
- 修改HomeScreen的跳转逻辑,等待页面返回后发送事件:
@override Widget build(BuildContext context) { return CupertinoPageScaffold( navigationBar: CupertinoNavigationBar( leading: const Icon(CupertinoIcons.line_horizontal_3), trailing: IconButton( onPressed: () async { // 等待页面返回 await Navigator.push( context, CupertinoPageRoute( builder: (BuildContext context) => const AddNewTransaction(), ), ); // 返回后重新加载交易数据 context.read<TransactionsBloc>().add(const TransactionsEvent.loadAll()); }, icon: const Icon(CupertinoIcons.add), ), ), child: const TransactionsList(), ); }
- 确保Bloc的
_loadTransactions方法每次都重新订阅流:
Future<void> _loadTransactions( TransactionsEvent event, Emitter<TransactionsState> emit) async { try { emit(state.copyWith(transactionStatus: TransactionsStatus.fetching)); // 每次加载都重新获取流实例 final transactionsStream = _transactionRepository.streamTransactions(); await emit.forEach( transactionsStream, onData: (transactions) { if (transactions.isEmpty) { return state.copyWith( transactionStatus: TransactionsStatus.initial, transactions: [], ); } return state.copyWith( transactionStatus: TransactionsStatus.fetchedSuccessfully, transactions: transactions, ); }, onError: (e, s) { return state.copyWith( transactionStatus: TransactionsStatus.fetchingFailed, message: 'Error loading transaction(s)', ); }, ); } catch (e) { add(const TransactionsEvent.error('Error loading transaction(s)')); } }
额外优化建议
- 确保
TransactionRepository是全局单例或在应用生命周期内唯一,避免重复创建导致BehaviorSubject实例被重置。 - 在修改交易的方法中,可直接更新
BehaviorSubject的数据源(而非重新查询全量数据),提升性能。
内容的提问来源于stack exchange,提问作者drunkZombie
相关产品推荐
相关产品推荐

