Flutter Firebase中如何等待Stream.listen执行完成后再返回结果
Flutter Stream问题
我需要通过geohash过滤数据库查询结果,因此使用了返回Stream的Geoflutterfire库。我尝试将Stream<List<>>中的每个DocumentSnapshot数据转换为List,功能本身可正常运行,但存在return语句在stream.listen()实际执行完成前就触发的问题。请问如何延迟return操作,待stream.listen处理完结果后再返回?我尝试使用await for(...)语法,但出现报错。
原始代码
import 'dart:async'; import 'package:cloud_firestore/cloud_firestore.dart'; import 'package:flutter_redux/flutter_redux.dart'; import 'package:geoflutterfire/geoflutterfire.dart'; import 'package:firebase_auth/firebase_auth.dart'; import 'package:firebase_storage/firebase_storage.dart'; import 'package:google_maps_flutter/google_maps_flutter.dart'; import 'package:uerto/models/index.dart'; class SearchApi { const SearchApi({required FirebaseAuth auth, required FirebaseFirestore firestore, required FirebaseStorage storage, required Geoflutterfire geo}) : _auth = auth, _firestore = firestore, _storage = storage, _geo = geo; final FirebaseAuth _auth; final FirebaseFirestore _firestore; final FirebaseStorage _storage; final Geoflutterfire _geo; Future<List<AppClient>> getClientList(LatLng location, String category, String subCategory, double radius, int limit) async{ final List<AppClient> newResult = <AppClient>[]; final GeoFirePoint center = _geo.point(latitude: location.latitude, longitude: location.longitude); final Query<Map<String, dynamic>> collectionReference = _firestore.collection('London$category/$subCategory/UID').limit(limit); const String field = 'position'; final Stream<List<DocumentSnapshot<Map<String, dynamic>>>> stream = _geo.collection(collectionRef: collectionReference).within(center: center, radius: radius, field: field); // ignore: always_specify_types stream.listen((List<DocumentSnapshot> documentList) { // ignore: always_specify_types, avoid_function_literals_in_foreach_calls documentList.forEach((DocumentSnapshot document) async { ///print(document.data()); final SearchUid searchUid = SearchUid.fromJson(document.data()); //print(searchUid.uid); final DocumentSnapshot<Map<String, dynamic>> client = await _firestore.collection('clients').doc(searchUid.uid).get(); final AppClient clientData = AppClient.fromJson(client.data()); print(clientData); newResult.add(clientData); }); }); await for(List<DocumentSnapshot> documentList in stream){ return newResult; } } }
报错截图

问题原因与解决方案
你代码的核心问题有三个:
- 对同一个单订阅Stream同时调用了
listen()和await for两次订阅,直接触发报错 listen回调中使用forEach包裹异步请求,forEach不会等待异步回调执行完成,导致数据还没加载完就执行返回- 流事件处理逻辑和返回逻辑拆分,无法保证执行顺序
修改后的完整代码
import 'dart:async'; import 'package:cloud_firestore/cloud_firestore.dart'; import 'package:flutter_redux/flutter_redux.dart'; import 'package:geoflutterfire/Geoflutterfire.dart'; import 'package:firebase_auth/firebase_auth.dart'; import 'package:firebase_storage/firebase_storage.dart'; import 'package:google_maps_flutter/google_maps_flutter.dart'; import 'package:uerto/models/index.dart'; class SearchApi { const SearchApi({required FirebaseAuth auth, required FirebaseFirestore firestore, required FirebaseStorage storage, required Geoflutterfire geo}) : _auth = auth, _firestore = firestore, _storage = storage, _geo = geo; final FirebaseAuth _auth; final FirebaseFirestore _firestore; final FirebaseStorage _storage; final Geoflutterfire _geo; Future<List<AppClient>> getClientList(LatLng location, String category, String subCategory, double radius, int limit) async{ final GeoFirePoint center = _geo.point(latitude: location.latitude, longitude: location.longitude); final Query<Map<String, dynamic>> collectionReference = _firestore.collection('London$category/$subCategory/UID').limit(limit); const String field = 'position'; final Stream<List<DocumentSnapshot<Map<String, dynamic>>>> stream = _geo.collection(collectionRef: collectionReference).within(center: center, radius: radius, field: field); // 只取流的第一波结果,不需要后续实时更新的话用first即可 final List<DocumentSnapshot<Map<String, dynamic>>> documentList = await stream.first; // 用Future.wait等待所有客户端文档查询完成 final List<AppClient> newResult = await Future.wait(documentList.map((document) async { final SearchUid searchUid = SearchUid.fromJson(document.data()!); final DocumentSnapshot<Map<String, dynamic>> client = await _firestore.collection('clients').doc(searchUid.uid).get(); return AppClient.fromJson(client.data()!); })); return newResult; } }
关键修改说明
- 删除了多余的
stream.listen调用,避免重复订阅报错,如果你需要后续实时监听位置更新返回流,可将方法返回值改为Stream<List<AppClient>>,使用asyncMap转换流数据即可 - 替换
forEach为Future.wait批量处理异步请求,确保所有客户端数据查询完成后再组装结果 - 直接通过
await stream.first获取第一波查询结果,不需要使用await for循环,逻辑更简洁 - 新增了
data()!的非空判断,避免空安全报错,你可根据自己的业务逻辑调整为空判断处理
如果需要持续监听位置更新实时返回最新的客户端列表,可以使用返回Stream的版本:
Stream<List<AppClient>> getClientListStream(LatLng location, String category, String subCategory, double radius, int limit) { final GeoFirePoint center = _geo.point(latitude: location.latitude, longitude: location.longitude); final Query<Map<String, dynamic>> collectionReference = _firestore.collection('London$category/$subCategory/UID').limit(limit); const String field = 'position'; return _geo.collection(collectionRef: collectionReference).within(center: center, radius: radius, field: field) .asyncMap((documentList) async { return await Future.wait(documentList.map((document) async { final SearchUid searchUid = SearchUid.fromJson(document.data()!); final DocumentSnapshot<Map<String, dynamic>> client = await _firestore.collection('clients').doc(searchUid.uid).get(); return AppClient.fromJson(client.data()!); })); }); }
内容的提问来源于stack exchange,提问作者FLTY
相关产品推荐
相关产品推荐

