Dart Stream — 非同期データストリームとリアクティブパイプライン
Stream はデータの川である — 連続的、到着した順に処理、すべてが一度に到着するのを待つ必要なし。
1. 学べること
- Stream の概念:単一購読 vs ブロードキャスト
- Stream の作成:fromIterable / periodic / イベント変換
- 演算子:map / where / expand / take / skip / distinct
- StreamController と StreamSubscription
- Bob のシナリオ:DataPipeline のリアルタイム注文ストリーム処理
2. 開発者のリアルな物語
(1) 課題:百万件のデータレコードを一度にメモリに読み込めない
Bob の DataPipeline は当初、EC 注文を処理する際に全 120 万件の注文をメモリに読み込んで処理していたため、メモリピークが 4GB に達し、OOM エラーでサーバーが 3 回クラッシュした。さらに、リアルタイムの注文が到着し続けるため、バッチ処理モデルは新しく到着したデータを扱えなかった。
(2) Stream による解決
Stream は水流のようにデータを 1 つずつ処理できる:到着した順に処理し、メモリ使用量を 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[データソース] --> B["StreamController"]
B --> C["where() フィルタ"]
C --> D["map() 変換"]
D --> E["expand() フラット化"]
E --> F["take() 制限"]
F --> G["listen() 消費"]
subgraph Stream Pipeline
C
D
E
F
end
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
| 種類 | 購読者 | リプレイ | ユースケース |
|---|---|---|---|
| 単一購読 Stream | 1 | リプレイ不可 | ファイル読み込み、HTTP レスポンス |
| ブロードキャスト Stream | 複数 | リプレイ不可 | UI イベント、リアルタイムデータ |
▶ サンプル:Stream の基礎
import 'dart:async';
void main() async {
// iterable から stream を作成
final stream = Stream.fromIterable([1, 2, 3, 4, 5]);
// stream を聞く
await for (final value in stream) {
print('Received: $value');
}
// 代替:コールバックで聞く
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'),
);
// stream の完了を待つ
await stream2.drain();
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
4. Stream の作成
▶ サンプル:fromIterable 作成
import 'dart:async';
void main() async {
// iterable から - 各要素を発行
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 バージョンにより結果が多少異なる場合があります。
▶ サンプル:periodic 作成
import 'dart:async';
void main() async {
// periodic stream - 一定間隔で値を発行
final ticker = Stream.periodic(
const Duration(seconds: 1),
(count) => 'Tick ${count + 1}',
);
// 最初の 5 チックのみ取得
await for (final tick in ticker.take(5)) {
print(tick);
}
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
▶ サンプル:StreamController 作成
import 'dart:async';
void main() async {
final controller = StreamController<String>();
// stream にデータを追加
controller.add('ORD-001');
controller.add('ORD-002');
controller.add('ORD-003');
// 完了時に stream を閉じる
controller.close();
// 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() |
なし | 空の ストリーム |
| 値 | Stream.value(v) |
単一値 | 単一値をラップ |
5. Stream 演算子
▶ サンプル:where — フィルタリング
import 'dart:async';
void main() async {
final amounts = Stream.fromIterable([1500.0, 50.0, 3200.0, 890.0, -100.0]);
// 100 以上の正の値のみフィルタ
final highValue = amounts.where((a) => a >= 100);
await for (final amount in highValue) {
print('\$${amount.toStringAsFixed(2)} USD');
}
// 出力: $1500.00 USD, $3200.00 USD, $890.00 USD
}
> 出力: ローカルの 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]);
// 金額をフォーマット済み文字列に変換
final formatted = amounts.map((a) => '\$${a.toStringAsFixed(2)} USD');
await for (final text in formatted) {
print(text);
}
// 出力: $1500.00 USD, $3200.00 USD, $890.00 USD
}
> 出力: ローカルの 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: 各注文 → 複数のアイテム
final items = orders.expand((order) =>
(order['items'] as List).cast<String>());
await for (final item in items) {
print('Item: $item');
}
// 出力: Item: Laptop, Item: Mouse, Item: Keyboard
}
> 出力: ローカルの 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]);
// 最初の 3 要素を取得
print('take(3):');
await data.take(3).forEach(print); // 1, 1, 2
// 最初の 2 をスキップ
print('skip(2):');
await Stream.fromIterable([1, 1, 2, 2, 3]).skip(2).forEach(print); // 2, 2, 3
// 連続する重複を削除
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
▶ サンプル:完全な 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();
// stream を購読
final subscription = orderStream.stream.listen(
(order) => print('Processing: ${order['id']}'),
onError: (e) => print('Error: $e'),
onDone: () => print('Stream closed'),
);
// 注文を追加
orderStream.addOrder({'id': 'ORD-001', 'amount': 1500.0});
orderStream.addOrder({'id': 'ORD-002', 'amount': 3200.0});
orderStream.addOrder({'id': 'ORD-003', 'amount': 890.0});
// stream を閉じる
orderStream.close();
// stream の終了を待つ
await subscription.asFuture();
}
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
(2) ブロードキャスト StreamController
▶ サンプル:ブロードキャストモード
import 'dart:async';
void main() async {
// ブロードキャストコントローラー - 複数の購読者
final controller = StreamController<String>.broadcast();
// 複数の購読者
final sub1 = controller.stream.listen((e) => print('Sub1: $e'));
final sub2 = controller.stream.listen((e) => print('Sub2: $e'));
// 両方の購読者が同じイベントを受信
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 のシナリオ:リアルタイム注文ストリーム処理
▶ サンプル:注文を処理する 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 {
// リアルタイム注文ストリームをシミュレート
final controller = StreamController<Order>();
// 処理パイプラインを構築
final revenueByCategory = <String, double>{};
var totalOrders = 0;
var completedOrders = 0;
controller.stream
.where((order) => order.amount > 0) // 無効をフィルタ
.where((order) => order.status == 'completed') // 完了のみ
.map((order) => order) // ここで変換可能
.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');
},
);
// 到着する注文をシミュレート
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')); // 無効
controller.add(Order(id: 'ORD-003', amount: 3200.0, status: 'pending', category: 'Electronics')); // 未完了
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 処理パイプライン
// 演算子付きのリアルタイム注文ストリーム
// ============================================
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) // 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();
// リアルタイム注文ストリームをシミュレート
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'), // 重複
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'),
];
// 注文を 1 つずつ発行(リアルタイムをシミュレート)
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 を複数回購読できますか? A: いいえ。単一購読 Stream は 1 回しか購読できず、2 回目の listen は StateError をスローします。複数の購読者が必要な場合は、ブロードキャスト Stream(
asBroadcastStream()またはStreamController.broadcast()経由)を使用してください。
Q: Stream 演算子は遅延的ですか? A: はい。map/where などの演算子はすぐには実行されず、購読されたときにのみ処理を開始します。これは「コールドストリーム」と呼ばれます。
Q: Stream でエラーを処理するにはどうすればよいですか? A: 3 つの方法があります:listen の onError コールバック、handleError 演算子、または
await forを try-catch でラップすることです。listen の onError コールバックの使用を推奨します。
Q: StreamController は閉じる必要がありますか? A: はい。controller.close() を呼び出して ストリーム を閉じ、購読者に ストリーム が終了したことを通知します。閉じないと購読者は無限に待機することになります。
Q: ブロードキャスト Stream の制限は何ですか? A: ブロードキャスト Stream は一時停止/再開をサポートせず、イベントが発行された後に購読した購読者は以前のイベントを逃します。「リアルタイム購読」シナリオに適しています。
Q: Stream はバックプレッシャーをサポートしますか? A: 単一購読 Stream は自然にバックプレッシャーをサポートします — 購読者(コンシューマ)の処理が遅い場合、生産者は待機します。ブロードキャスト Stream はバックプレッシャーをサポートしません。
📖 まとめ
- Stream は非同期イベントのシーケンスであり、単一購読(1 購読者)とブロードキャスト(複数購読者)に分かれる
- 作成方法:fromIterable(既存データ)、periodic(時間生成)、StreamController(手動制御)
- 演算子チェーン:where フィルタ → map 変換 → expand フラット化 → take/skip 制限 → distinct 重複排除
- StreamController は Stream の生産側、StreamSubscription は消費側
- Stream パイプラインにより百万件のデータを 1 つずつ処理でき、メモリ使用量が安定し、リアルタイム処理をサポートする
📝 練習問題
- 基礎(難易度 ⭐):
Stream.fromIterableを使って 10 件の金額を含む ストリーム を作成し、whereで 500 を超えるものをフィルタし、mapで USD 形式に変換し、listenで結果を出力してください。 - 中級(難易度 ⭐⭐):
StreamControllerを使ってリアルタイム注文ストリームをシミュレートし、5 秒間にわたって 1 秒ごとに 1 件の注文を追加してください。完了済み注文をフィルタし、カテゴリ別に売上を集約し、最終的にレポートを出力するパイプラインを構築します。 - 挑戦(難易度 ⭐⭐⭐):2 つの注文ストリームを 1 つに結合する「ストリームマージ」関数を実装してください(時間でインターリーブ)、重複排除して処理します。ヒント:
StreamGroup(asyncパッケージから)を使うか、StreamControllerで手動マージします。