Dart에서 스트림 만들기

Dart에서 스트림 만들기 (Creating streams in Dart)

스트림(stream)은 결과들의 시퀀스예요. 이 글에서는 스트림을 사용하는 법이 아니라, 직접 스트림을 만드는 방법을 여러 가지로 살펴볼게요.

출처: Creating streams in Dart

본문

  • 작성: Lasse Nielsen
  • 2013년 4월 (2021년 5월 업데이트)

dart:async 라이브러리는 많은 Dart API에 중요한 두 가지 타입을 담고 있어요: StreamFuture. Future가 단일 계산의 결과를 나타낸다면, stream은 결과의 시퀀스예요. stream을 듣고(listen) 있으면 결과들(데이터와 오류 모두)과 stream이 종료되는 순간을 통보받아요. 듣는 동안 멈췄다가(pause) 다시 재개할 수도 있고, stream이 완료되기 전에 듣기를 중단할 수도 있어요.

하지만 이 글은 stream을 사용하는 법이 아니라 직접 stream을 만드는 법에 관한 거예요. stream을 만드는 방법은 몇 가지가 있어요:

  • 기존 stream을 변환하기.
  • async* 함수로 처음부터 stream 만들기.
  • StreamController로 stream 만들기.

이 글은 각 방법의 코드를 보여 주고, stream을 올바르게 구현하는 데 도움이 되는 팁을 주는 내용이에요.

stream 사용에 대한 도움말은 비동기 프로그래밍: Streams를 참고하세요.

기존 stream 변환하기 (Transforming an existing stream)

stream을 만들 때 가장 흔한 경우는 이미 stream이 있고, 그 원본 stream의 이벤트를 바탕으로 새 stream을 만들고 싶을 때예요. 예를 들어 바이트 stream이 있어서 그 입력을 UTF-8 디코딩해 문자열 stream으로 바꾸고 싶을 수 있어요. 가장 일반적인 접근은 원본 stream의 이벤트를 기다렸다가 새 이벤트를 출력하는 새 stream을 만드는 거예요. 예시:

/// Splits a stream of consecutive strings into lines.
///
/// The input string is provided in smaller chunks through
/// the `source` stream.
Stream<String> lines(Stream<String> source) async* {
  // Stores any partial line from the previous chunk.
  var partial = '';
  // Wait until a new chunk is available, then process it.
  await for (final chunk in source) {
    var lines = chunk.split('\n');
    lines[0] = partial + lines[0]; // Prepend partial line.
    partial = lines.removeLast(); // Remove new partial line.
    for (final line in lines) {
      yield line; // Add lines to output stream.
    }
  }
  // Add final partial line to output stream, if any.
  if (partial.isNotEmpty) yield partial;
}

많은 흔한 변환에는 Stream이 제공하는 변환 메서드인 map(), where(), expand(), take() 등을 쓰면 돼요.

예를 들어 매 초마다 증가하는 카운터를 내보내는 stream인 counterStream이 있다고 가정해 볼게요. 다음과 같이 구현할 수 있어요:

var counterStream = Stream<int>.periodic(
  const Duration(seconds: 1),
  (x) => x,
).take(15);

이벤트를 빠르게 확인하고 싶다면 이런 코드를 쓰면 돼요:

counterStream.forEach(print); // Print an integer every second, 15 times.

stream 이벤트를 변환하려면 stream을 듣기 전에 map() 같은 변환 메서드를 호출하면 돼요. 그 메서드는 새 stream을 반환해요.

// Double the integer in each event.
var doubleCounterStream = counterStream.map((int x) => x * 2);
doubleCounterStream.forEach(print);

map() 대신 다음처럼 다른 변환 메서드를 써도 돼요:

.where((int x) => x.isEven) // Retain only even integer events.
.expand((x) => [x, x]) // Duplicate each event.
.take(5) // Stop after the first five events.

종종 변환 메서드 하나면 충분해요. 하지만 변환을 더 세밀하게 제어해야 한다면, Stream의 transform() 메서드에 StreamTransformer를 지정할 수 있어요. 플랫폼 라이브러리들은 많은 흔한 작업을 위한 stream transformer를 제공해요. 예를 들어 다음 코드는 dart:convert 라이브러리가 제공하는 utf8.decoderLineSplitter transformer를 사용해요.

Stream<List<int>> content = File('someFile.txt').openRead();
List<String> lines = await content
    .transform(utf8.decoder)
    .transform(const LineSplitter())
    .toList();

처음부터 stream 만들기 (Creating a stream from scratch)

새 stream을 만드는 한 가지 방법은 비동기 생성기(async*) 함수를 쓰는 거예요. 함수가 호출되면 stream이 만들어지고, stream이 들리기 시작하면 함수 본문이 실행되기 시작해요. 함수가 반환하면 stream이 닫혀요. 함수가 반환하기 전까지는 yieldyield* 문으로 stream에 이벤트를 내보낼 수 있어요.

일정한 간격으로 숫자를 내보내는 기본 예시를 볼게요:

Stream<int> timedCounter(Duration interval, [int? maxCount]) async* {
  int i = 0;
  while (true) {
    await Future.delayed(interval);
    yield i++;
    if (i == maxCount) break;
  }
}

이 함수는 Stream을 반환해요. 그 stream이 들리기 시작하면 본문이 실행되기 시작해요. 요청된 간격만큼 반복해서 지연한 뒤 다음 숫자를 yield해요. maxCount 매개변수를 생략하면 루프에 중단 조건이 없으므로, stream은 (리스너가 구독을 취소하기 전까지) 계속해서 점점 더 큰 숫자를 출력해요.

리스너가 취소하면(listen() 메서드가 반환한 StreamSubscription 객체의 cancel()을 호출해서) 다음에 본문이 yield 문에 도달할 때, 그 yield는 return 문처럼 동작해요. 둘러싼 finally 블록이 실행되고 함수가 종료돼요. 함수가 종료되기 전에 값을 yield하려 하면 그것은 실패하고 return처럼 동작해요.

함수가 마침내 종료되면 cancel() 메서드가 반환한 future가 완료돼요. 함수가 오류로 종료되면 그 future는 그 오류로 완료되고, 그렇지 않으면 null로 완료돼요.

또 하나 더 유용한 예시는 futures의 시퀀스를 stream으로 변환하는 함수예요:

Stream<T> streamFromFutures<T>(Iterable<Future<T>> futures) async* {
  for (final future in futures) {
    var result = await future;
    yield result;
  }
}

이 함수는 futures iterable에 새 future를 요청하고, 그 future를 기다렸다가 결과 값을 내보내고, 그리고 루프를 돌아요. future가 오류로 완료되면 stream도 그 오류로 완료돼요.

async* 함수가 아무것도 없는 곳에서 stream을 만들 일은 드물어요. 어딘가에서 데이터를 얻어야 하는데, 대부분 그 '어딘가'는 또 다른 stream이에요. 위의 futures 시퀀스 같은 경우처럼 데이터가 다른 비동기 이벤트 소스에서 오는 경우도 있지만, 많은 경우 async* 함수는 여러 데이터 소스를 다루기엔 너무 단순해요. 그때 StreamController 클래스가 등장해요.

StreamController 사용하기 (Using a StreamController)

stream의 이벤트가 async 함수가 순회할 수 있는 stream이나 futures가 아니라 프로그램의 여러 부분에서 온다면, StreamController를 사용해 stream을 만들고 채워 넣어요.

StreamController는 새 stream과, 언제 어디서든 stream에 이벤트를 추가할 수 있는 방법을 제공해요. 그 stream에는 리스너와 일시정지를 다루는 데 필요한 모든 로직이 들어 있어요. stream을 반환하고 controller는 여러분이 갖고 있으면 돼요.

다음 예시(stream_controller_bad.dart에서 가져온 것)는 이전 예시들의 timedCounter() 함수를 구현하기 위한, 기본적이지만 결함이 있는 StreamController 사용법을 보여 줘요. 이 코드는 반환할 stream을 만들고, futures도 stream 이벤트도 아닌 타이머 이벤트를 기반으로 데이터를 그 stream에 흘려 보내요.

// NOTE: This implementation is FLAWED!
// It starts before it has subscribers, and it doesn't implement pause.
Stream<int> timedCounter(Duration interval, [int? maxCount]) {
  var controller = StreamController<int>();
  int counter = 0;
  void tick(Timer timer) {
    counter++;
    controller.add(counter); // Ask stream to send counter values as event.
    if (maxCount != null && counter >= maxCount) {
      timer.cancel();
      controller.close(); // Ask stream to shut down and tell listeners.
    }
  }

  Timer.periodic(interval, tick); // BAD: Starts before it has subscribers.
  return controller.stream;
}

전처럼 timedCounter()가 반환한 stream을 이렇게 사용할 수 있어요:

var counterStream = timedCounter(const Duration(seconds: 1), 15);
counterStream.listen(print); // Print an integer every second, 15 times.

timedCounter() 구현에는 몇 가지 문제가 있어요:

  • 구독자가 생기기 전에 이벤트를 만들기 시작해요.
  • 구독자가 일시정지를 요청해도 계속 이벤트를 만들어 내요.

다음 섹션들이 보여 주듯이 StreamController를 만들 때 onListen, onPause 같은 콜백을 지정하면 이 두 문제를 모두 고칠 수 있어요.

구독 기다리기 (Waiting for a subscription)

원칙적으로 stream은 작업을 시작하기 전에 구독자를 기다려야 해요. async* 함수는 자동으로 그렇게 하지만, StreamController를 쓸 때는 완전히 통제권이 여러분에게 있어서 하지 말아야 할 때에도 이벤트를 추가할 수 있게 돼요. stream에 구독자가 없으면 StreamController는 이벤트를 버퍼링하는데, stream이 구독자를 영영 얻지 못하면 메모리 누수가 될 수 있어요.

stream을 사용하는 코드를 다음처럼 바꿔서 실행해 보세요:

void listenAfterDelay() async {
  var counterStream = timedCounter(const Duration(seconds: 1), 15);
  await Future.delayed(const Duration(seconds: 5));

  // After 5 seconds, add a listener.
  await for (final n in counterStream) {
    print(n); // Print an integer every second, 15 times.
  }
}

이 코드가 실행되면 처음 5초 동안은 stream이 작업을 하고 있는데도 아무것도 출력되지 않아요. 그 다음 리스너가 추가되고, 처음 5개 정도의 이벤트가 (StreamController가 버퍼링했으므로) 한꺼번에 출력돼요.

구독을 통보받으려면 StreamController를 만들 때 onListen 인자를 지정해요. stream이 첫 구독자를 얻으면 onListen 콜백이 호출돼요. onCancel 콜백을 지정하면 controller가 마지막 구독자를 잃을 때 호출돼요. 앞의 예시에서 Timer.periodic()은 다음 섹션에서 보듯 onListen 핸들러로 옮겨야 해요.

일시정지 상태 존중하기 (Honoring the pause state)

리스너가 일시정지를 요청했을 때 이벤트를 만들어 내는 건 피해야 해요. async* 함수는 stream 구독이 일시정지된 동안 yield 문에서 자동으로 멈춰요. 반면 StreamController는 일시정지 동안 이벤트를 버퍼링해요. 이벤트를 제공하는 코드가 일시정지를 존중하지 않으면 버퍼 크기가 무한정 커질 수 있어요. 게다가 리스너가 일시정지 직후 들기를 중단하면 버퍼를 만드는 데 쓴 작업이 낭비돼요.

일시정지 지원이 없으면 어떤 일이 일어나는지 보려면, stream을 사용하는 코드를 다음처럼 바꿔서 실행해 보세요:

void listenWithPause() {
  var counterStream = timedCounter(const Duration(seconds: 1), 15);
  late StreamSubscription<int> subscription;

  subscription = counterStream.listen((int counter) {
    print(counter); // Print an integer every second.
    if (counter == 5) {
      // After 5 ticks, pause for five seconds, then resume.
      subscription.pause(Future.delayed(const Duration(seconds: 5)));
    }
  });
}

5초의 일시정지가 끝나면 그동안 발생한 이벤트들이 한꺼번에 모두 수신돼요. stream의 소스가 일시정지를 존중하지 않고 계속 stream에 이벤트를 추가하기 때문이에요. 그래서 stream이 이벤트를 버퍼링하다가 일시정지가 풀리면 버퍼를 비우는 거예요.

stream_controller.dart에서 가져온 다음 timedCounter() 버전은 StreamController의 onListen, onPause, onResume, onCancel 콜백을 사용해 일시정지를 구현해요.

Stream<int> timedCounter(Duration interval, [int? maxCount]) {
  late StreamController<int> controller;
  Timer? timer;
  int counter = 0;

  void tick(_) {
    counter++;
    controller.add(counter); // Ask stream to send counter values as event.
    if (counter == maxCount) {
      timer?.cancel();
      controller.close(); // Ask stream to shut down and tell listeners.
    }
  }

  void startTimer() {
    timer = Timer.periodic(interval, tick);
  }

  void stopTimer() {
    timer?.cancel();
    timer = null;
  }

  controller = StreamController<int>(
    onListen: startTimer,
    onPause: stopTimer,
    onResume: startTimer,
    onCancel: stopTimer,
  );

  return controller.stream;
}

이 코드를 위의 listenWithPause() 함수와 함께 실행해 보세요. 일시정지 동안엔 카운터가 멈추고, 이후엔 잘 재개되는 것을 볼 수 있을 거예요.

일시정지 상태의 변화를 통보받으려면 onListen, onCancel, onPause, onResume 리스너를 모두 사용해야 해요. 구독 상태와 일시정지 상태가 동시에 변하면 onListen 또는 onCancel 콜백만 호출되기 때문이에요.

마지막 힌트들 (Final hints)

async* 함수를 사용하지 않고 stream을 만들 때 다음 팁들을 기억하세요:

  • 동기 controller를 쓸 때 조심해요 — 예를 들어 StreamController(sync: true)로 만든 것을 말해요. 일시정지되지 않은 동기 controller에 이벤트를 보내면(예: EventSink가 정의한 add(), addError(), close() 메서드 사용) 그 이벤트는 즉시 stream의 모든 리스너에게 전달돼요. stream 리스너는 리스너를 추가한 코드가 완전히 반환할 때까지 절대 호출되면 안 되는데, 동기 controller를 잘못된 시점에 쓰면 이 약속을 깨뜨려 좋은 코드조차 실패하게 만들 수 있어요. 동기 controller는 피하세요.
  • StreamController를 쓴다면, onListen 콜백이 listen 호출이 StreamSubscription을 반환하기 전에 호출된다는 점을 기억하세요. onListen 콜백이 구독이 이미 존재한다는 것에 의존하지 않게 하세요. 예를 들어 다음 코드에서는 subscription 변수에 유효한 값이 생기기 전에 onListen 이벤트가 발생해요(그리고 handler가 호출돼요).
    subscription = stream.listen(handler);
    
  • StreamController가 정의한 onListen, onPause, onResume, onCancel 콜백은 stream의 리스너 상태가 변할 때 stream에 의해 호출되지만, 이벤트가 발생하는 중이나 다른 상태 변경 핸들러를 호출하는 중에는 절대 호출되지 않아요. 그런 경우 상태 변경 콜백은 이전 콜백이 끝날 때까지 지연돼요.
  • Stream 인터페이스를 직접 구현하려고 하지 마세요. 이벤트, 콜백, 리스너 추가·제거 사이의 상호 작용을 미묘하게 잘못 잡기 쉽거든요. 새 stream의 listen 호출을 구현할 때는 항상 기존 stream(가능하면 StreamController에서 온 것)을 사용하세요.
  • Stream 클래스를 상속해 더 많은 기능을 가진 클래스를 만들 수 있긴 해요(listen 메서드와 그 위의 추가 기능을 구현해서). 하지만 이는 사용자들이 고려해야 할 새 타입을 도입하기 때문에 일반적으로 권장하지 않아요. "Stream이면서 더 많은 것"인 클래스 대신, 종종 "Stream을 가지면서 더 많은 것"인 클래스를 만들 수 있어요.

더 알아보기