开篇:从"拉取"到"推送"的思维转变
传统编程是"拉取式"的:调用函数、等待返回、处理结果。而 Flutter 应用充满的是"推送式"事件——用户点击、网络响应、传感器数据、消息推送。Dart 的 Stream 正是为这种异步事件流而生:它把"一连串随时间到达的事件"抽象为一个可监听、可转换、可组合的数据流。
从底层 StreamController 到响应式库 RxDart,从 StreamBuilder 驱动 UI 到与 BLoC/Riverpod 架构结合,Stream 贯穿了 Flutter 异步编程的方方面面。本章将从零开始构建完整的响应式数据流知识体系,并给出搜索防抖、实时数据、多数据源合并等实战案例。
一、Stream 基础与 StreamController
1.1 Stream 是什么
Stream 是一系列异步事件的序列。它有三种关键事件:
- 数据事件:
onData—— 正常数据到达 - 错误事件:
onError—— 出错了(可继续) - 完成事件:
onDone—— 流结束(不再有事件)
Stream<int> countStream() async* {
for (var i = 1; i <= 3; i++) {
yield i; // 产生一个数据事件
await Future.delayed(const Duration(seconds: 1));
}
}
void main() {
final sub = countStream().listen(
(data) => print('收到: $data'), // onData
onError: (e) => print('错误: $e'), // onError
onDone: () => print('流结束'), // onDone
);
// sub.cancel() 可取消订阅
}
1.2 StreamController 的手动控制
StreamController 允许你手动向流中推入事件,是连接"外部世界"与"流世界"的桥梁:
import 'dart:async';
class ChatService {
final _controller = StreamController<String>();
Stream<String> get messages => _controller.stream; // 暴露只读流
void send(String message) {
if (!_controller.isClosed) {
_controller.add(message); // 推入数据
}
}
void close() {
_controller.close(); // 触发 onDone
}
}
| StreamController 成员 | 作用 |
|---|---|
add(data) | 推送数据事件 |
addError(error) | 推送错误事件 |
addStream(other) | 把另一个流接进来 |
close() | 结束流并触发 onDone |
stream | 对外暴露的可监听流 |
一句话:
Stream是"事件的时间轴",StreamController是"轴上按钮"——add 推事件,listen 收事件,cancel 关订阅。
二、单订阅流 vs 广播流
2.1 两种流的本质区别
| 特性 | 单订阅流(Single-subscription) | 广播流(Broadcast) |
|---|---|---|
| 监听者数量 | 只能有 1 个 | 可以有很多个 |
| 事件缓存 | 未监听时事件会排队 | 未监听时事件直接丢弃 |
| 数据源 | 一次性数据源(文件、HTTP) | 多接收方共享(点击、通知) |
| 创建方式 | StreamController() | StreamController.broadcast() |
// 单订阅流:重复 listen 会报错
final single = StreamController<int>();
final sub1 = single.stream.listen((v) => print(v));
// final sub2 = single.stream.listen(...); // ❌ StateError: Stream has already been listened
// 广播流:多个监听者共享
final broadcast = StreamController<int>.broadcast();
broadcast.stream.listen((v) => print('A: $v'));
broadcast.stream.listen((v) => print('B: $v'));
broadcast.add(1); // 两个监听者都会收到
2.2 如何选择
// 场景 1:HTTP 响应、文件读取 → 单订阅流
Future<Stream<String>> readFile(String path) {
return File(path).openRead(); // 天然单订阅
}
// 场景 2:全局事件总线 → 广播流
class AppBus {
AppBus._();
static final _controller = StreamController<AppEvent>.broadcast();
static Stream<AppEvent> get events => _controller.stream;
static void emit(AppEvent event) => _controller.add(event);
}
一句话:默认用单订阅流(更省内存、语义更清晰),只有"一个事件要被多个监听者接收"时才用广播流。
三、StreamBuilder 与事件驱动 UI
3.1 StreamBuilder 的核心用法
StreamBuilder 监听一个 Stream,事件到达时重建 UI,是"响应式 UI"的最直接体现:
import 'package:flutter/material.dart';
class CounterStreamPage extends StatefulWidget {
const CounterStreamPage({super.key});
@override
State<CounterStreamPage> createState() => _CounterStreamPageState();
}
class _CounterStreamPageState extends State<CounterStreamPage> {
final _controller = StreamController<int>();
int _count = 0;
void _increment() => _controller.add(++_count);
@override
void dispose() {
_controller.close(); // 关键:关闭控制器释放资源
super.dispose();
}
@override
Widget build(BuildContext context) {
return Scaffold(
appBar: AppBar(title: const Text('Stream 计数器')),
body: Center(
child: StreamBuilder<int>(
stream: _controller.stream,
builder: (context, snapshot) {
return Text(
'Count: ${snapshot.data ?? 0}',
style: const TextStyle(fontSize: 32),
);
},
),
),
floatingActionButton: FloatingActionButton(
onPressed: _increment,
child: const Icon(Icons.add),
),
);
}
}
3.2 StreamBuilder 的 snapshot 状态
| snapshot 状态 | 含义 | 典型 UI |
|---|---|---|
ConnectionState.none | 尚未开始监听 | 初始占位 |
ConnectionState.waiting | 等待第一个事件 | 加载中 |
ConnectionState.active | 已收到事件,流未结束 | 展示数据 |
ConnectionState.done | 流已结束 | 无更多数据提示 |
snapshot.hasError | 流中有错误 | 错误提示 + 重试 |
StreamBuilder<List<Post>>(
stream: postService.fetchPosts(),
builder: (context, snapshot) {
if (snapshot.hasError) return const Text('加载失败');
if (!snapshot.hasData) return const CircularProgressIndicator();
final posts = snapshot.data!;
return ListView.builder(
itemCount: posts.length,
itemBuilder: (context, i) => ListTile(title: Text(posts[i].title)),
);
},
)
一句话:
StreamBuilder把"流事件"直接映射为"UI 状态"——这正是响应式编程在 UI 层的落地:数据变化,界面自动跟随。
四、RxDart 操作符
4.1 为什么需要 RxDart
原生 Stream 缺少丰富的转换操作符。RxDart 在 Dart Stream API 之上提供了类 ReactiveX 的操作符组合能力。
4.2 三大高频操作符
import 'package:rxdart/rxdart.dart';
// 1. debounce:防抖——只在停止输入 500ms 后才发出
final searchStream = searchController.stream
.debounceTime(const Duration(milliseconds: 500));
// 2. combineLatest:合并多个流的最新值
final validForm = Rx.combineLatest2<String, String, bool>(
emailController.stream,
passwordController.stream,
(email, password) =>
email.isNotEmpty && password.isNotEmpty,
);
// 3. switchMap:切换到最新的事件流(丢弃旧的)
final results = queryStream.switchMap((query) {
return searchService.search(query); // 每个 query 产生新流
});
4.3 常用操作符速查
| 操作符 | 作用 | 典型场景 |
|---|---|---|
map | 转换每个事件 | 数据模型映射 |
where | 过滤事件 | 只处理合法输入 |
debounceTime | 防抖 | 搜索框、滑动结束 |
combineLatest | 合并最新值 | 表单多字段校验 |
switchMap | 切换到最新流 | 竞态(快速切换 tab) |
distinct | 去重连续相同值 | 避免重复请求 |
retry | 出错重试 | 网络抖动 |
// 组合使用:搜索框完整响应式链路
Stream<List<Product>> searchStream(Stream<String> queryStream) {
return queryStream
.debounceTime(const Duration(milliseconds: 400))
.distinct() // 相同 query 不重复请求
.switchMap((q) => api.searchProducts(q)) // 只取最新结果
.retry(2); // 失败重试
}
一句话:RxDart 的威力在于"组合"——debounce + distinct + switchMap 一条链,就解决了搜索框 80% 的复杂度。
五、流与 BLoC / Riverpod 结合
5.1 BLoC:Stream 是主角
BLoC 模式的本质就是"UI 发事件,BLoC 用 Stream 回状态":
import 'package:flutter_bloc/flutter_bloc.dart';
sealed class SearchEvent {}
class SearchChanged extends SearchEvent {
final String query;
SearchChanged(this.query);
}
class SearchBloc extends Bloc<SearchEvent, List<String>> {
final ProductApi api;
SearchBloc(this.api) : super(const []) {
on<SearchChanged>(_onSearchChanged);
}
Future<void> _onSearchChanged(
SearchChanged event,
Emitter<List<String>> emit,
) async {
// 注意:BLoC 内部可用 Stream + 防抖,或借助 bloc 的 transformEvents
final results = await api.search(event.query);
emit(results);
}
}
5.2 BLoC 内嵌 RxDart 防抖
BLoC 官方支持用 transformEvents 注入操作符,实现流的防抖:
class SearchBloc extends Bloc<SearchEvent, SearchState> {
SearchBloc(this.api) : super(const SearchState()) {
on<SearchChanged>(_search, transformer: debounce(const Duration(milliseconds: 400)));
}
@override
Stream<SearchEvent> transformEvents(
Stream<SearchEvent> events,
Stream<SearchState> Function(SearchEvent event) next,
) {
// 用 RxDart 的 debounceTime 做全局防抖
return events.debounceTime(const Duration(milliseconds: 400)).switchMap(next);
}
}
5.3 StreamProvider:Riverpod 中的流
Riverpod 的 StreamProvider 直接拥抱 Stream:
import 'package:flutter_riverpod/flutter_riverpod.dart';
final messagesProvider = StreamProvider<List<Message>>((ref) {
final chatService = ref.watch(chatServiceProvider);
return chatService.messages; // 返回 Stream
});
class ChatView extends ConsumerWidget {
const ChatView({super.key});
@override
Widget build(BuildContext context, WidgetRef ref) {
final asyncMessages = ref.watch(messagesProvider);
return asyncMessages.when(
data: (messages) => ListView.builder(
itemCount: messages.length,
itemBuilder: (context, i) => ListTile(title: Text(messages[i].body)),
),
loading: () => const CircularProgressIndicator(),
error: (e, _) => Text('连接失败: $e'),
);
}
}
一句话:BLoC 把 Stream 当架构主轴,Riverpod 的 StreamProvider 则把流"声明式"地接入组件——无论哪种,Stream 都是两者与异步世界对话的通用语言。
六、背压与资源清理
6.1 背压(Backpressure)
当"生产者"产出事件的速率超过"消费者"处理速率时,就产生了背压。Dart 的 Stream 默认策略:
| 策略 | 行为 | 适用场景 |
|---|---|---|
| 缓存 | 事件进入缓冲区排队 | 单订阅流默认,速率可控 |
| 丢弃 | 来不及处理的直接丢弃 | 传感器高频数据 |
| 合并 | 用最新值覆盖旧值 | 进度条、实时位置 |
// 通过操作符实现"丢弃中间值,只取最新"
final latestValue = highFrequencyStream
.sampleTime(const Duration(milliseconds: 100)) // 每 100ms 取最新
.distinct();
6.2 资源清理的三件事
流相关的资源如果不清理,会引发内存泄漏:
class _StreamPageState extends State<StreamPage> {
StreamSubscription<Data>? _sub;
final _controller = StreamController<Data>();
@override
void initState() {
super.initState();
_sub = _controller.stream.listen(_handleData); // 1. 订阅要记录
}
void _handleData(Data data) => setState(() {});
@override
void dispose() {
_sub?.cancel(); // 2. 取消订阅
_controller.close(); // 3. 关闭控制器
super.dispose();
}
}
| 资源 | 清理方式 | 后果(未清理) |
|---|---|---|
StreamSubscription | cancel() | 回调继续触发,内存泄漏 |
StreamController | close() | 事件阻塞、资源不释放 |
Timer / Ticker | cancel() / dispose() | 定时回调泄漏 |
一句话:背压用"采样/丢弃/合并"策略应对高频事件,资源清理用"订阅 cancel + 控制器 close"双保险——两者都是流编程的必修课。
七、错误处理与重试
7.1 流的错误传递
Stream 的 onError 与 Future 的 catchError 类似,但流可能携带多个错误:
stream.listen(
(data) => print(data),
onError: (Object e, StackTrace st) {
debugPrint('错误: $e');
// 单订阅流中 onError 后流可能继续,广播流中错误事件同样广播
},
)
7.2 结合 Future 的异步转换
在流中执行异步操作(网络请求)时,用 asyncMap 或 switchMap:
// asyncMap:每个事件都触发异步请求,结果依序输出
final enriched = userIds.asyncMap((id) => api.fetchUser(id));
// switchMap:竞态安全,只保留最后一次请求的结果
final latestUser = userIds.switchMap((id) => api.fetchUser(id));
7.3 重试与超时
import 'package:rxdart/rxdart.dart';
// 失败自动重试 3 次,间隔递增
final retryStream = api.fetchData().retryWhen((errors) {
var attempt = 0;
return errors.map((e) {
attempt++;
if (attempt >= 3) throw e; // 超过 3 次抛给上层
return e;
}).delay(const Duration(seconds: 1));
});
// 超时:5 秒无事件则触发错误
final timeoutStream = api.fetchData().timeout(const Duration(seconds: 5));
一句话:流的错误是"事件的一部分"——用
asyncMap/switchMap做异步转换,用retryWhen/timeout做健壮性兜底。
八、实战案例:搜索防抖 + 竞态控制
把本章知识串起来,实现一个完整的搜索页:
import 'dart:async';
import 'package:flutter/material.dart';
import 'package:rxdart/rxdart.dart';
class SearchPage extends StatefulWidget {
const SearchPage({super.key});
@override
State<SearchPage> createState() => _SearchPageState();
}
class _SearchPageState extends State<SearchPage> {
final _queryController = StreamController<String>();
late final Stream<List<String>> _results;
final _fakeApi = FakeApi();
@override
void initState() {
super.initState();
_results = _queryController.stream
.debounceTime(const Duration(milliseconds: 400)) // 防抖
.distinct() // 相同输入不重复请求
.switchMap(_fakeApi.search) // 竞态安全
.onErrorReturnWith((_) => const []); // 错误兜底
}
@override
void dispose() {
_queryController.close();
super.dispose();
}
@override
Widget build(BuildContext context) {
return Scaffold(
appBar: AppBar(title: const Text('响应式搜索')),
body: Column(
children: [
Padding(
padding: const EdgeInsets.all(16),
child: TextField(
onChanged: _queryController.add, // 输入即推入流
decoration: const InputDecoration(
hintText: '输入关键词',
prefixIcon: Icon(Icons.search),
),
),
),
Expanded(
child: StreamBuilder<List<String>>(
stream: _results,
builder: (context, snapshot) {
if (snapshot.hasError) return const Center(child: Text('出错了'));
if (!snapshot.hasData) {
return const Center(child: CircularProgressIndicator());
}
final results = snapshot.data!;
if (results.isEmpty) {
return const Center(child: Text('暂无结果'));
}
return ListView.builder(
itemCount: results.length,
itemBuilder: (context, i) => ListTile(
leading: const Icon(Icons.article),
title: Text(results[i]),
),
);
},
),
),
],
),
);
}
}
class FakeApi {
Future<List<String>> search(String query) async {
await Future.delayed(const Duration(milliseconds: 300));
return List.generate(10, (i) => '$query 结果 $i');
}
}
实战要点回顾
| 关注点 | 处理方式 |
|---|---|
| 输入防抖 | debounceTime(400ms) |
| 重复请求 | distinct() |
| 竞态 | switchMap(只保留最新) |
| 错误 | onErrorReturnWith |
| 资源 | dispose 中 close() 控制器 |
一句话:响应式数据流把"状态变化"变成"数据流的变化"——搜索、聊天、实时位置等场景,用一条操作符链即可优雅地表达完整逻辑。
九、总结
Stream 与 RxDart 构成了 Flutter 响应式编程的完整工具箱:
| 能力 | 工具 | 关键要点 |
|---|---|---|
| 事件建模 | Stream / StreamController | 单订阅 vs 广播 |
| UI 绑定 | StreamBuilder | snapshot 四状态 |
| 组合转换 | RxDart 操作符 | debounce/switchMap/combineLatest |
| 架构融合 | BLoC / Riverpod | 流是异步的统一语言 |
| 健壮性 | retry / timeout | 错误是流的一部分 |
| 资源管理 | cancel + close | 双保险防泄漏 |
一句话:从"回调金字塔"到"流水线式操作符链",Stream 让异步事件变得可声明、可组合、可预测——掌握它,你就掌握了 Flutter 异步世界的核心语法。
相关阅读
- https://plumephp.com/flutter-async-networking/ — 异步编程与网络通信
- https://plumephp.com/flutter-state-management/ — 状态管理(BLoC 与 Riverpod 的流用法)
- https://plumephp.com/flutter-form-validation/ — 表单验证(输入流防抖校验)
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。