Dart Stream — 非同期データストリームとリアクティブパイプライン

Stream はデータの川である — 連続的、到着した順に処理、すべてが一度に到着するのを待つ必要なし。

1. 学べること


2. 開発者のリアルな物語

(1) 課題:百万件のデータレコードを一度にメモリに読み込めない

Bob の DataPipeline は当初、EC 注文を処理する際に全 120 万件の注文をメモリに読み込んで処理していたため、メモリピークが 4GB に達し、OOM エラーでサーバーが 3 回クラッシュした。さらに、リアルタイムの注文が到着し続けるため、バッチ処理モデルは新しく到着したデータを扱えなかった。

(2) Stream による解決

Stream は水流のようにデータを 1 つずつ処理できる:到着した順に処理し、メモリ使用量を 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[データソース] --> 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
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。
種類 購読者 リプレイ ユースケース
単一購読 Stream 1 リプレイ不可 ファイル読み込み、HTTP レスポンス
ブロードキャスト Stream 複数 リプレイ不可 UI イベント、リアルタイムデータ

▶ サンプル:Stream の基礎

DART
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();
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

4. Stream の作成

▶ サンプル:fromIterable 作成

DART
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');
  }
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

▶ サンプル:periodic 作成

DART
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);
  }
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

▶ サンプル:StreamController 作成

DART
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');
  }
}
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() なし 空の ストリーム
Stream.value(v) 単一値 単一値をラップ

5. Stream 演算子

▶ サンプル:where — フィルタリング

DART
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
}
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]);

  // 金額をフォーマット済み文字列に変換
  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
}
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: 各注文 → 複数のアイテム
  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
}
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]);

  // 最初の 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
}
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

▶ サンプル:完全な 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();

  // 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();
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

(2) ブロードキャスト StreamController

▶ サンプル:ブロードキャストモード

DART
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));
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

7. Bob のシナリオ:リアルタイム注文ストリーム処理

▶ サンプル:注文を処理する 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 {
  // リアルタイム注文ストリームをシミュレート
  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));
}
TEXT
> 出力: ローカルの DartPad または `dart run` で実行してください。この Dart コースの全例は Dart 3.x / Flutter 3.x ベースです。SDK バージョンにより結果が多少異なる場合があります。

8. 完全なサンプル:DataPipeline Stream パイプライン

DART
// ============================================
// 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();
}
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 を複数回購読できますか? 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 はバックプレッシャーをサポートしません。


📖 まとめ


📝 練習問題

  1. 基礎(難易度 ⭐)Stream.fromIterable を使って 10 件の金額を含む ストリーム を作成し、where で 500 を超えるものをフィルタし、map で USD 形式に変換し、listen で結果を出力してください。
  2. 中級(難易度 ⭐⭐)StreamController を使ってリアルタイム注文ストリームをシミュレートし、5 秒間にわたって 1 秒ごとに 1 件の注文を追加してください。完了済み注文をフィルタし、カテゴリ別に売上を集約し、最終的にレポートを出力するパイプラインを構築します。
  3. 挑戦(難易度 ⭐⭐⭐):2 つの注文ストリームを 1 つに結合する「ストリームマージ」関数を実装してください(時間でインターリーブ)、重複排除して処理します。ヒント:StreamGroupasync パッケージから)を使うか、StreamController で手動マージします。

← 前のレッスン | 次のレッスン →

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%