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:

ScenariosExample
User input eventsButton clicks, mouse movement, keyboard input
File readingRead large files line by line
Network dataWebSocket messages, real-time APIs
TimerTimer events that fire once per second
State changesState management streams in Flutter.

Example

Basic usage of Stream — listening with listen:

void main() {
  // 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:

MethodsFunction
subscription.pause()Pause receiving data
subscription.resume()Resume receiving data
subscription.cancel()Cancel subscription
subscription.isPausedWhether 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

// Generate a finite data stream
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

import 'dart:async';

// 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

// async* marks this as an async generator function
// 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

void main() {
  // 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

import 'dart:async';

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:

FeaturesSingle-Subscription StreamBroadcast Stream
Number of subscribersOnly oneMultiple
Data replayConsume from the beginningOnly receives data after subscribing
Typical scenariosFile reading, HTTP responsesButton clicks, status notifications
Creation MethodStreamController()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

void main() async {
  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
  3
other extensions