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

如何在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 an Observable<QuerySnapshot<MyObject>> that emits updates whenever the DataStore's data for MyObject changes (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 uses await for to listen to the stream, and yields the snapshot.items (a List<MyObject>) every time a new snapshot is received. This creates the exact Stream<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 onError callback 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 18:27:47