Dart: Dart Stream — 异步数据流与响应式管道

Stream 是数据的河流 — 源源不断,边来边处理,不用等全部到齐。

1. 你将学到


2. 一个开发者的真实故事

(1) 痛点:百万级数据无法一次性加载到内存

Bob 的 DataPipeline 处理电商订单时,最初把全部 1,200,000 条订单加载到内存再处理,导致内存峰值 4GB,服务器 OOM 崩溃了 3 次。而且实时订单不断到来,批处理模式无法处理新到达的数据。

(2) Stream 的解法

Stream 让数据像水流一样逐条处理:来一条处理一条,内存占用稳定在 50MB。Stream 操作符(where/map/reduce)让处理管道声明式定义。

DART
orderStream
  .where((order) => order.status == 'completed')
  .map((order) => order.amount)
  .fold(0, (sum, amount) => sum + amount)
  .then((total) => print('Revenue: \$$total USD'));
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

(3) 收益


3. Stream 基础

(1) 单订阅 vs 广播

100%
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
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。
类型 订阅者数量 重放 适用场景
单订阅 Stream 1 个 不可重放 文件读取、HTTP 响应
广播 Stream 多个 不可重放 UI 事件、实时数据

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:Stream 基础

DART
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();
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

4. Stream 创建

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:fromIterable 创建

DART
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');
  }
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:periodic 创建

DART
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);
  }
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:StreamController 创建

DART
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');
  }
}
TEXT 📖 仅展示
> **输出:** 在本地 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 操作符

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:where — 过滤

DART
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
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:map — 转换

DART
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
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:expand — 展开

DART
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
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:take / skip / distinct

DART
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
}
TEXT 📖 仅展示
> **输出:** 在本地 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

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:StreamController 完整用法

DART
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();
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

(2) 广播 StreamController

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:广播模式

DART
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));
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

7. Bob 场景:实时订单流处理

▶ 示例

TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

:Stream 管道处理订单

DART
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));
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

8. 完整示例:DataPipeline Stream 管道

DART
// ============================================
// 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();
}
TEXT 📖 仅展示
> **输出:** 在本地 DartPad 或 `dart run` 执行。Dart 课程所有示例基于 Dart 3.x / Flutter 3.x,运行结果会因 SDK 版本略有差异。

输出:

TEXT 📖 仅展示
=== 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 不支持背压。


📖 小节


📝 作业

  1. 基础题(难度⭐):用 Stream.fromIterable 创建一个包含 10 个金额的流,用 where 过滤出大于 500 的,用 map 转换为 USD 格式,用 listen 打印结果。
  2. 进阶题(难度⭐⭐):用 StreamController 实现一个模拟的实时订单流,每秒添加 1 个订单,持续 5 秒。构建管道过滤已完成订单,按类别聚合收入,最后输出报告。
  3. 挑战题(难度⭐⭐⭐):实现一个"合并流"功能,将两个订单流合并为一个流(按时间交错),去重后处理。提示:用 StreamGroup(async 包)或 StreamController 手动合并。

← 上一课 | 下一课 →

Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏