Futureが「1つの値を将来返す非同期処理」であるのに対し、Stream は「時間の経過とともに複数の値やイベントが次々と流れてくる非同期のデータシーケンス」です。
WebSocketのメッセージ受信、ファイルの逐次読み込み、ボタンの連続タップイベント、センサー情報のリアルタイム受信、そしてFlutterの状態管理(BLoCパターンやRxDartなど)において Stream は極めて重要な役割を果たします。
この章では、Stream の基本概念、await for による購読、StreamController、そして非同期ジェネレータ(async* / yield)を学びます。
Stream の概念(Push型のデータフロー)Stream<T> は、非同期に次々と発生するデータ(またはエラー)のパイプラインです。
購読者(リスナー)はデータが準備できたタイミングで通知を受け取ります(Push型)。
+--------------------------------------------------------+ | Stream<int> | | ---[ 1 ]-----[ 2 ]------[ 3 ]------| (完了 / Done) | | 1秒後 2秒後 3秒後 | +--------------------------------------------------------+
Streamは以下の3つの通知をリスナーへ送信します。
T)。DartのStreamには2つのモードが存在します。
listen)できます。ファイルの読み込みやHTTPレスポンスなどに適しています。stream.asBroadcastStream() で変換できます。await for によるStreamの購読async 関数内では、await for ループを使って Stream から送られてくるデータを同期ループのように1件ずつ順次処理できます。
ストリームが完了(Done)するまでループが継続します。
Stream<int> countStream(int max) async* {
for (int i = 1; i <= max; i++) {
yield i;
}
}
void main() async {
print('Stream受信開始');
// await for による順次受信
await for (final number in countStream(3)) {
print('受信データ: $number');
}
print('Stream完了');
}dart run await_for_demo.dartStream受信開始 受信データ: 1 受信データ: 2 受信データ: 3 Stream完了
Iterable と同様に、Stream に対しても where(フィルタリング)、map(変換)、take(指定件数取得)などのパイプライン処理をメソッドチェーンで適用できます。
void main() async {
final numbersStream = Stream.fromIterable([1, 2, 3, 4, 5, 6]);
// 偶数だけを抽出し、10倍に変換して最初の2件を取得
final transformedStream = numbersStream
.where((n) => n.isEven)
.map((n) => n * 10)
.take(2);
await for (final val in transformedStream) {
print('変換後データ: $val');
}
}dart run stream_operators.dart変換後データ: 20 変換後データ: 40
StreamController と StreamSubscriptionプログラム側から能動的にイベントを生成・送信(Push)したい場合は、dart:async ライブラリの StreamController<T> を使用します。
また、購読の開始と停止を制御するために StreamSubscription を管理します。
controller.sink.add(value) でデータを送信し、購読は stream.listen(...) で行います。
import 'dart:async';
void main() async {
// StreamController の作成
final controller = StreamController<String>();
// 購読を開始 (StreamSubscription を保持)
final subscription = controller.stream.listen(
(message) => print('受信: $message'),
onError: (error) => print('エラー受信: $error'),
onDone: () => print('ストリーム終了 (Done)'),
);
// イベントの発行 (sink経由)
controller.sink.add('第1通知: ユーザーログイン');
controller.sink.add('第2通知: メッセージ受信');
controller.sink.addError('第3通知: 接続不安定警告');
// ストリームを閉じる
await controller.close();
// 不要になったら購読を破棄(メモリリーク防止)
await subscription.cancel();
}dart run stream_controller_demo.dart受信: 第1通知: ユーザーログイン 受信: 第2通知: メッセージ受信 エラー受信: 第3通知: 接続不安定警告 ストリーム終了 (Done)
async* と yield(非同期ジェネレータ関数)async*(非同期ジェネレータ) を使うと、関数の内部から時間の経過とともに複数の値を yield で順次ストリームへ送り出すことができます。
Stream<String> tickStream(int count) async* {
for (int i = 1; i <= count; i++) {
await Future.delayed(const Duration(milliseconds: 50));
yield 'Tick #$i';
}
}
void main() async {
await for (final tick in tickStream(3)) {
print(tick);
}
}dart run async_generator.dartTick #1 Tick #2 Tick #3
この章では、Dartのリアクティブプログラミングを支える Stream について学びました。
Stream<T>: 連続する非同期イベントを流すデータパイプライン(Push型)。await for: ストリームから流れてくる要素をループで簡潔に順次受信。where, map, take などのオペレータでイベントシーケンスを加工。StreamController: 能動的にイベントを発行し、StreamSubscription で購読とキャンセルを管理。async* / yield: 非同期ジェネレータ関数で宣言的にストリームを生成。次の 第11章 では、実務において不可欠な エラーハンドリング(例外処理とResult型パターン) について学びます。
指定された秒数からゼロまでカウントダウンするストリーム関数を作成してください。
Stream<int> countdown(int from) を async* で定義する。from から 0 まで1ずつ減らしながら yield する(各ステップで Future.delayed(Duration(milliseconds: 50)) を待つ)。main() で countdown(3) を呼び出し、await for で受け取って '残り: X秒'、最後に 'カウントダウン終了!' と出力する。// ここに関数を定義してください
void main() async {
// ここで動作確認を行ってください
}dart run practice10_1.dartメッセージ通知システムを StreamController を使って構築してください。
class NotificationHub を作成する。
StreamController<String> を保持する。void sendNotification(String message) でイベントを送信する。Stream<String> get onNotification でストリームを公開する。Future<void> dispose() でコントローラを閉じる。main() でインスタンスを作成し、onNotification を listen して受信ログを出力する。dispose() を呼ぶ。import 'dart:async';
// ここにクラスを定義してください
void main() async {
// ここで動作確認を行ってください
}dart run practice10_2.dart