Dart Stream
Stream is the mechanism in Dart for handling sequences of continuous asynchronous events.
If Future is a one-time asynchronous result, then Stream is a repeated, continuous asynchronous data flow.
This chapter introduces the concept of Stream, await for listening, StreamController creation, and the differences between single-subscription and broadcast streams.
Streams and event sequences
A Stream is like a conveyor belt; data arrives one by one over time.
You don't need to wait for all the data at once; instead, you process each piece as it arrives.
Typical use cases for Streams:
| Scenarios | Example |
|---|---|
| User input events | Button clicks, mouse movement, keyboard input |
| File reading | Read large files line by line |
| Network data | WebSocket messages, real-time APIs |
| Timer | Timer events that fire once per second |
| State changes | State management streams in Flutter. |
Example
Basic usage of Stream — listening with listen:
// Create a Stream: emit a number every 1 second
var stream = Stream<int>.periodic(
Duration(seconds: 1),
(count) => count + 1, // count starts from 0
);
print('Start listening to Stream...');
// listen subscribes to the Stream
var subscription = stream.listen(
(data) {
// Called whenever new data arrives
print('EXAMPLE received data: $data');
},
onError: (error) {
// Called when Stream encounters an error
print('Error: $error');
},
onDone: () {
// Called when Stream is closed
print('Stream closed');
},
);
// Cancel subscription after 5 seconds
Future.delayed(Duration(seconds: 5), () {
subscription.cancel();
print('Subscription cancelled');
});
}
开始监听 Stream... EXAMPLE 收到数据: 1 EXAMPLE 收到数据: 2 EXAMPLE 收到数据: 3 EXAMPLE 收到数据: 4 EXAMPLE 收到数据: 5 已取消订阅
listen() returns a StreamSubscription object, which you can use to control the subscription:
| Methods | Function |
|---|---|
| subscription.pause() | Pause receiving data |
| subscription.resume() | Resume receiving data |
| subscription.cancel() | Cancel subscription |
| subscription.isPaused | Whether it is in a paused state |
await for listening
Besides the listen() callback approach, you can also use an await for loop to consume a Stream.
await for makes the logic of processing a Stream as clear as an ordinary for loop.
Example
Stream<int> countStream(int max) async* {
for (int i = 1; i <= max; i++) {
await Future.delayed(Duration(milliseconds: 300));
yield i; // yield emits data to the Stream
}
}
Future<void> main() async {
print('Start consuming Stream with await for...');
// await for: wait for each data item to arrive and process them one by one
await for (var value in countStream(5)) {
print('EXAMPLE count: $value');
}
print('Stream consumption complete');
}
开始用 await for 消费 Stream... EXAMPLE 计数: 1 EXAMPLE 计数: 2 EXAMPLE 计数: 3 EXAMPLE 计数: 4 EXAMPLE 计数: 5 Stream 消费完毕
Characteristics of await for:
- The syntax is similar to a normal for loop, but each iteration waits for the next piece of data to arrive.
- When the Stream closes, the loop ends automatically.
- It can only be used in async functions.
- If you need to exit early, use break (just like a normal for loop).
await for is suitable for scenarios where you need to process all data sequentially, such as reading all lines from a file. listen() is suitable for scenarios where you need to respond to events continuously, such as button clicks. The two can replace each other, but each has its own strengths.
StreamController creates Stream
StreamController is a tool for manually creating and controlling Streams.
You can add data, errors, or close it at any time.
Example
// Use StreamController to implement a simple countdown timer
class CountdownTimer {
final StreamController<int> _controller = StreamController<int>();
Timer? _timer;
int _remaining = 0;
// Expose Stream for external subscription
Stream<int> get tickStream => _controller.stream;
// Start countdown
void start(int seconds) {
_remaining = seconds;
// Immediately send the initial value
_controller.add(_remaining);
_timer = Timer.periodic(Duration(seconds: 1), (timer) {
_remaining--;
if (_remaining > 0) {
_controller.add(_remaining); // Send data
} else {
_controller.add(0); // Send the final 0
_controller.close(); // Close Stream
timer.cancel();
}
});
}
// Cancel the countdown
void cancel() {
_timer?.cancel();
_controller.addError('Countdown canceled'); // Send error
_controller.close();
}
// Release resources
void dispose() {
_controller.close();
}
}
Future<void> main() async {
var timer = CountdownTimer();
// Subscribe to the countdown event
timer.tickStream.listen(
(remaining) {
print('EXAMPLE countdown: $remaining seconds');
},
onError: (error) {
print('Error: $error');
},
onDone: () {
print('Countdown finished!');
},
);
timer.start(5);
// Wait for the countdown to complete
await Future.delayed(Duration(seconds: 6));
timer.dispose();
}
EXAMPLE 倒计时: 5 秒 EXAMPLE 倒计时: 4 秒 EXAMPLE 倒计时: 3 秒 EXAMPLE 倒计时: 2 秒 EXAMPLE 倒计时: 1 秒 EXAMPLE 倒计时: 0 秒 倒计时结束!
async* generator functions
If you only need to simply generate a sequence of data, async* is more convenient than StreamController.
Example
// The return type must be Stream
Stream<String> readLinesAsync() async* {
var lines = ['First line', 'Second line', 'Third line', 'EXAMPLE'];
for (var line in lines) {
await Future.delayed(Duration(milliseconds: 500));
yield line; // yield emits data to the Stream
}
// Stream automatically closes when the function ends
}
// Async generator with error handling
Stream<int> generateNumbersWithError() async* {
for (int i = 1; i <= 5; i++) {
await Future.delayed(Duration(milliseconds: 300));
if (i == 3) {
throw Exception('Number 3 had an error!');
}
yield i;
}
}
Future<void> main() async {
print('Reading line by line:');
await for (var line in readLinesAsync()) {
print(' $line');
}
print('\n'Generator with errors:');
try {
await for (var num in generateNumbersWithError()) {
print(' Number: $num');
}
} catch (e) {
print(' Caught error: $e');
}
}
逐行读取: 第一行 第二行 第三行 EXAMPLE 带错误的生成器: 数字: 1 数字: 2 捕获错误: Exception: 数字 3 出错了!
Single-subscription vs. broadcast streams
Dart Streams are divided into two types: single-subscription streams and broadcast streams.
Single-subscription Stream
This is the default Stream type.
It can only be subscribed to by one listener, suitable for scenarios where you consume it once from start to finish.
Example
// Single-subscription stream (default)
var stream = Stream.fromIterable([1, 2, 3]);
// First subscription: OK
stream.listen((data) => print('Subscriber 1: $data'));
// Second subscription: Error! A single-subscription stream cannot be subscribed to multiple times
// stream.listen((data) => print('Subscriber 2: $data')); // runtime error
}
Broadcast Stream
Broadcast streams allow multiple listeners to subscribe simultaneously, suitable for event broadcast scenarios.
Example
void main() {
// Create a broadcast stream
var controller = StreamController<int>.broadcast();
// Multiple subscribers can listen at the same time
controller.stream.listen(
(data) => print('EXAMPLE Subscriber A: received $data'),
);
controller.stream.listen(
(data) => print('EXAMPLE Subscriber B: received $data'),
);
// Send data, both subscribers will receive it
controller.add(1);
controller.add(2);
controller.add(3);
// A late subscription can still receive subsequent data (but not previous data)
Future.delayed(Duration(seconds: 1), () {
controller.stream.listen(
(data) => print('EXAMPLE Late Subscriber C: received $data'),
);
controller.add(4);
controller.close();
});
}
EXAMPLE 订阅者A: 收到 1 EXAMPLE 订阅者B: 收到 1 EXAMPLE 订阅者A: 收到 2 EXAMPLE 订阅者B: 收到 2 EXAMPLE 订阅者A: 收到 3 EXAMPLE 订阅者B: 收到 3 EXAMPLE 订阅者A: 收到 4 EXAMPLE 订阅者B: 收到 4 EXAMPLE 迟到的订阅者C: 收到 4
Comparison of the two types of Streams:
| Features | Single-Subscription Stream | Broadcast Stream |
|---|---|---|
| Number of subscribers | Only one | Multiple |
| Data replay | Consume from the beginning | Only receives data after subscribing |
| Typical scenarios | File reading, HTTP responses | Button clicks, status notifications |
| Creation Method | StreamController() | StreamController.broadcast() |
A "late subscriber" to a broadcast stream can only receive events after subscribing; previous events have already been missed. This is the most important behavioral difference between broadcast streams and single-subscription streams.
Common Stream transformation methods
Stream provides a series of methods for transforming and processing data, similar to the functional methods of List.
Example
var numbers = Stream.fromIterable([1, 2, 3, 4, 5, 6]);
// map: transform each data item
var doubled = numbers.map((n) => n * 2);
print('Doubled:');
await for (var n in doubled) {
print(' $n');
}
// where: filter data
var source = Stream.fromIterable([10, 15, 20, 25, 30]);
var even = source.where((n) => n % 2 == 0);
print('Even:');
await for (var n in even) {
print(' $n');
}
// take: take only the first N items
var infinite = Stream.periodic(
Duration(milliseconds: 100),
(i) => i + 1,
);
print('Take the first 3 only:');
await for (var n in infinite.take(3)) {
print(' EXAMPLE: $n');
}
// skip: skip the first N items
var data = Stream.fromIterable([1, 2, 3, 4, 5]);
print('Skip the first 2:');
await for (var n in data.skip(2)) {
print(' $n');
}
// distinct: remove duplicates
var duplicates = Stream.fromIterable([1, 2, 2, 3, 3, 3]);
print('Distinct:');
await for (var n in duplicates.distinct()) {
print(' $n');
}
}
翻倍: 2 4 6 8 10 12 偶数: 10 20 30 只取前 3 个: EXAMPLE: 1 EXAMPLE: 2 EXAMPLE: 3 跳过前 2 个: 3 4 5 去重: 1 2 3other extensions