Dart: Dart Stream — 异步数据流与响应式管道
Stream 是数据的河流 — 源源不断,边来边处理,不用等全部到齐。
1. 你将学到
- Stream 概念:单订阅 vs 广播(Single-subscription vs Broadcast)
- Stream 创建:fromIterable / periodic / event-transformed
- 操作符:map / where / expand / take / skip / distinct
- StreamController 与 StreamSubscription
- Bob 场景:DataPipeline 实时订单流处理
2. 一个开发者的真实故事
(1) 痛点:百万级数据无法一次性加载到内存
Bob 的 DataPipeline 处理电商订单时,最初把全部 1,200,000 条订单加载到内存再处理,导致内存峰值 4GB,服务器 OOM 崩溃了 3 次。而且实时订单不断到来,批处理模式无法处理新到达的数据。
(2) Stream 的解法
Stream 让数据像水流一样逐条处理:来一条处理一条,内存占用稳定在 50MB。Stream 操作符(where/map/reduce)让处理管道声明式定义。
orderStream
.where((order) => order.status == 'completed')
.map((order) => order.amount)
.fold(0, (sum, amount) => sum + amount)
.then((total) => print('Revenue: \$$total USD'));
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
(3) 收益
- 内存从 4GB 降到 50MB,不再 OOM
- 实时处理新到达的订单,延迟从"等批处理"降到"秒级"
- Stream 操作符链让数据处理逻辑声明式、可组合
3. Stream 基础
(1) 单订阅 vs 广播
flowchart LR
A[Data Source] --> B["StreamController"]
B --> C["where() Filter"]
C --> D["map() Transform"]
D --> E["expand() Flatten"]
E --> F["take() Limit"]
F --> G["listen() Consume"]
subgraph Stream Pipeline
C
D
E
F
end
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
| 类型 | 订阅者数量 | 重放 | 适用场景 |
|---|---|---|---|
| 单订阅 Stream | 1 个 | 不可重放 | 文件读取、HTTP 响应 |
| 广播 Stream | 多个 | 不可重放 | UI 事件、实时数据 |
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:Stream 基础
import 'dart:async';
void main() async {
// Create stream from iterable
final stream = Stream.fromIterable([1, 2, 3, 4, 5]);
// Listen to stream
await for (final value in stream) {
print('Received: $value');
}
// Alternative: listen with callback
final stream2 = Stream.fromIterable(['ORD-001', 'ORD-002', 'ORD-003']);
stream2.listen(
(order) => print('Order: $order'),
onDone: () => print('Stream completed'),
onError: (e) => print('Error: $e'),
);
// Wait for stream to complete
await stream2.drain();
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
4. Stream 创建
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:fromIterable 创建
import 'dart:async';
void main() async {
// From iterable - emits each element
final orders = Stream.fromIterable([
{'id': 'ORD-001', 'amount': 1500.0},
{'id': 'ORD-002', 'amount': 3200.0},
{'id': 'ORD-003', 'amount': 890.0},
]);
await for (final order in orders) {
print('Order: ${order['id']} - \$${order['amount']} USD');
}
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:periodic 创建
import 'dart:async';
void main() async {
// Periodic stream - emits values at regular intervals
final ticker = Stream.periodic(
const Duration(seconds: 1),
(count) => 'Tick ${count + 1}',
);
// Take only first 5 ticks
await for (final tick in ticker.take(5)) {
print(tick);
}
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:StreamController 创建
import 'dart:async';
void main() async {
final controller = StreamController``<String>``();
// Add data to the stream
controller.add('ORD-001');
controller.add('ORD-002');
controller.add('ORD-003');
// Close the stream when done
controller.close();
// Consume the stream
await for (final order in controller.stream) {
print('Processing: $order');
}
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
| 创建方式 | 语法 | 数据源 | 适用场景 |
|---|---|---|---|
| fromIterable | Stream.fromIterable(list) |
已有集合 | 测试、小数据集 |
| periodic | Stream.periodic(duration, fn) |
定时生成 | 定时器、轮询 |
| StreamController | StreamController``<T>``() |
手动添加 | 事件源、实时数据 |
| empty | Stream.empty() |
无 | 空流 |
| value | Stream.value(v) |
单值 | 包裹单值 |
5. Stream 操作符
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:where — 过滤
import 'dart:async';
void main() async {
final amounts = Stream.fromIterable([1500.0, 50.0, 3200.0, 890.0, -100.0]);
// Filter only positive amounts >= 100
final highValue = amounts.where((a) => a >= 100);
await for (final amount in highValue) {
print('\$${amount.toStringAsFixed(2)} USD');
}
// Output: $1500.00 USD, $3200.00 USD, $890.00 USD
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:map — 转换
import 'dart:async';
void main() async {
final amounts = Stream.fromIterable([1500.0, 3200.0, 890.0]);
// Transform amounts to formatted strings
final formatted = amounts.map((a) => '\$${a.toStringAsFixed(2)} USD');
await for (final text in formatted) {
print(text);
}
// Output: $1500.00 USD, $3200.00 USD, $890.00 USD
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:expand — 展开
import 'dart:async';
void main() async {
final orders = Stream.fromIterable([
{'id': 'ORD-001', 'items': ['Laptop', 'Mouse']},
{'id': 'ORD-002', 'items': ['Keyboard']},
]);
// Expand: each order → multiple items
final items = orders.expand((order) =>
(order['items'] as List).cast``<String>``());
await for (final item in items) {
print('Item: $item');
}
// Output: Item: Laptop, Item: Mouse, Item: Keyboard
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:take / skip / distinct
import 'dart:async';
void main() async {
final data = Stream.fromIterable([1, 1, 2, 2, 3, 3, 4, 5]);
// Take first 3 elements
print('take(3):');
await data.take(3).forEach(print); // 1, 1, 2
// Skip first 2
print('skip(2):');
await Stream.fromIterable([1, 1, 2, 2, 3]).skip(2).forEach(print); // 2, 2, 3
// Remove consecutive duplicates
print('distinct:');
await Stream.fromIterable([1, 1, 2, 2, 3, 3]).distinct().forEach(print); // 1, 2, 3
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
| 操作符 | 作用 | 返回类型 |
|---|---|---|
where |
过滤元素 | Stream<T> |
map |
转换元素 | Stream<R> |
expand |
一对多展开 | Stream<T> |
take |
取前 N 个 | Stream<T> |
skip |
跳过前 N 个 | Stream<T> |
distinct |
去连续重复 | Stream<T> |
takeWhile |
条件取 | Stream<T> |
skipWhile |
条件跳 | Stream<T> |
6. StreamController 与 StreamSubscription
(1) StreamController
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:StreamController 完整用法
import 'dart:async';
class OrderStream {
final _controller = StreamController<Map<String, dynamic>>();
Stream<Map<String, dynamic>> get stream => _controller.stream;
void addOrder(Map<String, dynamic> order) => _controller.add(order);
void addError(Object error) => _controller.addError(error);
void close() => _controller.close();
}
void main() async {
final orderStream = OrderStream();
// Subscribe to the stream
final subscription = orderStream.stream.listen(
(order) => print('Processing: ${order['id']}'),
onError: (e) => print('Error: $e'),
onDone: () => print('Stream closed'),
);
// Add orders
orderStream.addOrder({'id': 'ORD-001', 'amount': 1500.0});
orderStream.addOrder({'id': 'ORD-002', 'amount': 3200.0});
orderStream.addOrder({'id': 'ORD-003', 'amount': 890.0});
// Close the stream
orderStream.close();
// Wait for stream to finish
await subscription.asFuture();
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
(2) 广播 StreamController
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:广播模式
import 'dart:async';
void main() async {
// Broadcast controller - multiple listeners
final controller = StreamController``<String>``.broadcast();
// Multiple subscribers
final sub1 = controller.stream.listen((e) => print('Sub1: $e'));
final sub2 = controller.stream.listen((e) => print('Sub2: $e'));
// Both subscribers receive the same event
controller.add('ORD-001');
controller.add('ORD-002');
controller.close();
await Future.delayed(const Duration(milliseconds: 100));
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
7. Bob 场景:实时订单流处理
▶ 示例
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
:Stream 管道处理订单
import 'dart:async';
class Order {
final String id;
final double amount;
final String status;
final String category;
Order({required this.id, required this.amount, required this.status, required this.category});
@override
String toString() => 'Order($id, \$${amount.toStringAsFixed(2)}, $status, $category)';
}
void main() async {
// Simulate real-time order stream
final controller = StreamController``<Order>``();
// Build processing pipeline
final revenueByCategory = <String, double>{};
var totalOrders = 0;
var completedOrders = 0;
controller.stream
.where((order) => order.amount > 0) // Filter invalid
.where((order) => order.status == 'completed') // Only completed
.map((order) => order) // Could transform here
.listen(
(order) {
completedOrders++;
totalOrders++;
revenueByCategory.update(
order.category,
(v) => v + order.amount,
ifAbsent: () => order.amount,
);
print(' Processed: ${order.id} - \$${order.amount.toStringAsFixed(2)} USD');
},
onDone: () {
print('\n=== Stream Processing Complete ===');
print('Completed orders: $completedOrders');
for (final entry in revenueByCategory.entries) {
print(' ${entry.key}: \$${entry.value.toStringAsFixed(2)} USD');
}
final total = revenueByCategory.values.fold(0.0, (a, b) => a + b);
print('Total revenue: \$${total.toStringAsFixed(2)} USD');
},
);
// Simulate incoming orders
controller.add(Order(id: 'ORD-001', amount: 1500.0, status: 'completed', category: 'Electronics'));
controller.add(Order(id: 'ORD-002', amount: -50.0, status: 'completed', category: 'Books')); // Invalid
controller.add(Order(id: 'ORD-003', amount: 3200.0, status: 'pending', category: 'Electronics')); // Not completed
controller.add(Order(id: 'ORD-004', amount: 890.0, status: 'completed', category: 'Clothing'));
controller.add(Order(id: 'ORD-005', amount: 2100.0, status: 'completed', category: 'Electronics'));
controller.close();
await Future.delayed(const Duration(milliseconds: 100));
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
8. 完整示例:DataPipeline Stream 管道
// ============================================
// DataPipeline Stream Processing Pipeline
// Real-time order stream with operators
// ============================================
import 'dart:async';
class Order {
final String id;
final double amount;
final String status;
final String category;
final String region;
const Order({
required this.id,
required this.amount,
required this.status,
required this.category,
this.region = 'US',
});
double get taxAmount => amount * 0.08;
double get totalWithTax => amount + taxAmount;
@override
String toString() =>
'Order($id, \$${amount.toStringAsFixed(2)}, $status, $category, $region)';
}
class StreamPipeline {
final _controller = StreamController``<Order>``();
final Map<String, double> _revenue = {};
final Map<String, int> _count = {};
int _totalProcessed = 0;
int _totalSkipped = 0;
late StreamSubscription``<Order>`` _subscription;
Stream``<Order>`` get stream => _controller.stream;
StreamPipeline() {
_subscription = _controller.stream
.where((o) => o.amount > 0)
.where((o) => o.status != 'cancelled')
.distinct((o) => o.id) // Deduplicate by ID
.listen(
_processOrder,
onError: (e) => print('Pipeline error: $e'),
onDone: _printSummary,
);
}
void _processOrder(Order order) {
_totalProcessed++;
_revenue.update(
order.category,
(v) => v + order.totalWithTax,
ifAbsent: () => order.totalWithTax,
);
_count.update(order.category, (v) => v + 1, ifAbsent: () => 1);
}
void addOrder(Order order) => _controller.add(order);
void addError(Object error) => _controller.addError(error);
Future``<void>`` close() async {
await _controller.close();
await _subscription.asFuture();
}
void _printSummary() {
final totalRevenue = _revenue.values.fold(0.0, (a, b) => a + b);
final totalCount = _count.values.fold(0, (a, b) => a + b);
print('\n=== DataPipeline Stream Report ===');
print('Processed: $_totalProcessed orders');
print('Skipped: $_totalSkipped records');
print('Revenue: \$${totalRevenue.toStringAsFixed(2)} USD');
print('Orders: $totalCount');
print('\n--- By Category ---');
final sorted = _revenue.entries.toList()
..sort((a, b) => b.value.compareTo(a.value));
for (final entry in sorted) {
final count = _count[entry.key] ?? 0;
final avg = entry.value / count;
print(' ${entry.key}:');
print(' Orders: $count');
print(' Revenue: \$${entry.value.toStringAsFixed(2)} USD');
print(' Average: \$${avg.toStringAsFixed(2)} USD');
}
}
}
void main() async {
final pipeline = StreamPipeline();
// Simulate real-time order stream
final orders = [
Order(id: 'ORD-001', amount: 1500.0, status: 'completed', category: 'Electronics'),
Order(id: 'ORD-002', amount: -50.0, status: 'completed', category: 'Books'),
Order(id: 'ORD-003', amount: 3200.0, status: 'pending', category: 'Electronics'),
Order(id: 'ORD-004', amount: 890.0, status: 'completed', category: 'Clothing'),
Order(id: 'ORD-001', amount: 1500.0, status: 'completed', category: 'Electronics'), // Duplicate
Order(id: 'ORD-005', amount: 2100.0, status: 'completed', category: 'Electronics'),
Order(id: 'ORD-006', amount: 500.0, status: 'cancelled', category: 'Books'),
Order(id: 'ORD-007', amount: 120.0, status: 'completed', category: 'Books'),
];
// Emit orders one by one (simulating real-time)
for (final order in orders) {
pipeline.addOrder(order);
}
await pipeline.close();
}
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
输出:
=== DataPipeline Stream Report ===
Processed: 5 orders
Skipped: 0 records
Revenue: $8154.00 USD
Orders: 5
--- By Category ---
Electronics:
Orders: 3
Revenue: $7236.00 USD
Average: $2412.00 USD
Clothing:
Orders: 1
Revenue: $961.20 USD
Average: $961.20 USD
Books:
Orders: 1
Revenue: $129.60 USD
Average: $129.60 USD
❓ 常见问题
Q:Stream 和 Future 有什么区别? A:Future 代表一个异步结果(0 或 1 个值),Stream 代表异步事件序列(0 到多个值)。Stream 是 Future 的推广 — 多值版本。
Q:单订阅 Stream 可以多次 listen 吗? A:不可以。单订阅 Stream 只能被 listen 一次,第二次会抛 StateError。需要多个监听者时用广播 Stream(
asBroadcastStream()或StreamController.broadcast())。
Q:Stream 操作符是惰性的吗? A:是的。map/where 等操作符不会立即执行,只有在 listen 时才开始处理。这叫"冷流"(cold stream)。
Q:如何处理 Stream 中的错误? A:三种方式:listen 的 onError 回调、handleError 操作符、try-catch 包裹 await for。推荐用 listen 的 onError。
Q:StreamController 用完需要关闭吗? A:是的。调用 controller.close() 关闭流,通知监听者流结束。不关闭会导致监听者永远等待。
Q:广播 Stream 有什么限制? A:广播 Stream 不提供暂停/恢复功能,且如果监听者在事件发出后才订阅,会错过之前的事件。适合"实时订阅"场景。
Q:Stream 可以背压(backpressure)吗? A:单订阅 Stream 天然支持背压 — listen 消费速度慢时,生产者会等待。广播 Stream 不支持背压。
📖 小节
- Stream 是异步事件序列,分单订阅(1 个监听者)和广播(多个监听者)
- 创建方式:fromIterable(已有数据)、periodic(定时生成)、StreamController(手动控制)
- 操作符链:where 过滤 → map 转换 → expand 展开 → take/skip 限量 → distinct 去重
- StreamController 是 Stream 的生产端,StreamSubscription 是消费端
- Stream 管道让百万级数据逐条处理,内存稳定,支持实时处理
📝 作业
- 基础题(难度⭐):用 Stream.fromIterable 创建一个包含 10 个金额的流,用 where 过滤出大于 500 的,用 map 转换为 USD 格式,用 listen 打印结果。
- 进阶题(难度⭐⭐):用 StreamController 实现一个模拟的实时订单流,每秒添加 1 个订单,持续 5 秒。构建管道过滤已完成订单,按类别聚合收入,最后输出报告。
- 挑战题(难度⭐⭐⭐):实现一个"合并流"功能,将两个订单流合并为一个流(按时间交错),去重后处理。提示:用 StreamGroup(async 包)或 StreamController 手动合并。