如何在Repository类中将Amplify.DataStore.observeQuery转换为可在Bloc中使用、返回类型为Stream<List<MyObject>>的async* Stream
Convert Amplify.DataStore.observeQuery to Stream<List> for Bloc (async* Implementation)
Great question! Converting Amplify's observeQuery output to a Dart Stream<List<MyObject>> that works seamlessly with Bloc is straightforward once you understand how to bridge Amplify's Observable with Dart's Stream API. Here's a step-by-step implementation:
Step 1: Repository Method Implementation
In your repository class, create an async* generator function that converts the Observable returned by observeQuery into a stream of List<MyObject>:
import 'package:amplify_flutter/amplify_flutter.dart'; import 'package:your_app/models/MyObject.dart'; class MyRepository { // Basic implementation for observing all MyObject entries Stream<List<MyObject>> observeMyObjects() async* { // Convert Amplify's Observable to a standard Dart Stream final queryStream = Amplify.DataStore.observeQuery(MyObject.classType).asStream(); // Listen for incoming snapshots and yield the updated list each time await for (final QuerySnapshot<MyObject> snapshot in queryStream) { yield snapshot.items; } } // Optional: Filtered/sorted version for more specific use cases Stream<List<MyObject>> observeFilteredMyObjects(String searchTerm) async* { final queryStream = Amplify.DataStore.observeQuery( MyObject.classType, where: MyObject.TITLE.contains(searchTerm), // Example filter logic sortBy: [MyObject.UPDATED_AT.descending()], // Example sort order ).asStream(); await for (final snapshot in queryStream) { yield snapshot.items; } } }
How This Works:
Amplify.DataStore.observeQuery()returns anObservable<QuerySnapshot<MyObject>>that emits updates whenever the DataStore's data forMyObjectchanges (create, update, delete).- We use
.asStream()to convert this Observable into a standard Dart Stream, which plays nicely with Dart's async patterns. - The
async*generator function usesawait forto listen to the stream, andyields thesnapshot.items(aList<MyObject>) every time a new snapshot is received. This creates the exactStream<List<MyObject>>you need.
Step 2: Using the Stream in Bloc
Now you can subscribe to this stream in your Bloc to emit states whenever the data updates. Make sure to manage the subscription to avoid memory leaks:
import 'package:bloc/bloc.dart'; import 'package:your_app/repositories/my_repository.dart'; import 'package:your_app/models/MyObject.dart'; // Events abstract class MyEvent {} class StartObservingMyObjects extends MyEvent {} class UpdateSearchTerm extends MyEvent { final String searchTerm; UpdateSearchTerm(this.searchTerm); } // States abstract class MyState {} class MyObjectsLoading extends MyState {} class MyObjectsLoaded extends MyState { final List<MyObject> objects; MyObjectsLoaded(this.objects); } class MyObjectsError extends MyState { final String message; MyObjectsError(this.message); } class MyBloc extends Bloc<MyEvent, MyState> { final MyRepository _repository; StreamSubscription<List<MyObject>>? _subscription; MyBloc(this._repository) : super(MyObjectsLoading()) { on<StartObservingMyObjects>((event, emit) { _subscribeToObjects(emit); }); on<UpdateSearchTerm>((event, emit) { _subscribeToObjects(emit, searchTerm: event.searchTerm); }); } void _subscribeToObjects(Emitter<MyState> emit, {String searchTerm = ""}) { // Cancel existing subscription first to avoid duplicate stream emissions _subscription?.cancel(); // Pick the appropriate stream based on search term final stream = searchTerm.isEmpty ? _repository.observeMyObjects() : _repository.observeFilteredMyObjects(searchTerm); _subscription = stream.listen( (objects) => emit(MyObjectsLoaded(objects)), onError: (error) => emit(MyObjectsError(error.toString())), ); } @override Future<void> close() { // Cancel subscription when Bloc is disposed to prevent memory leaks _subscription?.cancel(); return super.close(); } }
Key Notes for Bloc Integration:
- Always cancel existing subscriptions before creating a new one (especially when updating filters/sort order) to avoid multiple streams emitting conflicting states.
- Handle errors in the stream's
onErrorcallback to emit meaningful error states to your UI. - Override the
close()method to cancel the subscription when the Bloc is disposed—this is critical for preventing memory leaks.
Why This Is Ideal for Bloc:
- The
async*generator function creates a clean, readable stream that aligns with Dart's asynchronous patterns. - Every DataStore change triggers a new state emission in the Bloc, keeping your UI in perfect sync with the latest data.
- You can easily extend the repository method to include filters, sorting, or pagination as your app's needs grow.
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

