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

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方法)已关闭。

需要解决:

  1. 在StatelessWidget架构下,如何维持流不关闭?
  2. 若无法维持,返回交易列表页时重新订阅流的推荐方案是什么?

解决方案

方案一:修正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事件重新订阅流:

  1. 修改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(),
  );
}
  1. 确保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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 20:22:15