스트림 사용하기: 단일 구독 스트림과 브로드캐스트 스트림

스트림 사용하기: 단일 구독 스트림과 브로드캐스트 스트림

스트림은 비동기 이벤트의 연속이에요. 요청하는 순간 다음 이벤트를 받는 일반 Iterable과 달리, 준비된 이벤트가 생기면 스트림이 알려주는 형태라고 보면 돼요. Dart의 비동기 프로그래밍은 FutureStream 두 클래스로 특징지어집니다. 이 글에서는 스트림 이벤트를 받는 방법, 에러 이벤트, 그리고 두 종류의 스트림을 어떻게 나누고 각각을 어떻게 다루는지 살펴볼게요.

출처: Dart 공식 문서 — Asynchronous programming: Streams

스트림 이벤트 받기

스트림은 다양한 방법으로 만들 수 있지만(이건 다른 주제예요) 사용 방식은 모두 같아요. 비동기 for 루프(보통 그냥 await for라고 불러요)가 for 루프가 Iterable을 순회하듯 스트림의 이벤트를 순회합니다.

Future<int> sumStream(Stream<int> stream) async {
  var sum = 0;
  await for (final value in stream) {
    sum += value;
  }
  return sum;
}

이 코드는 정수 이벤트 스트림의 각 이벤트를 받아 더한 뒤, (Future로 감싼) 합계를 반환해요. 루프 본문이 끝나면 함수는 다음 이벤트가 오거나 스트림이 끝날 때까지 멈춥니다. 함수에 async 키워드가 표시되어 있어야 await for 루프를 쓸 수 있어요.

아래 예시는 async* 함수로 정수 스트림을 만들어 위 코드를 시험해 봅니다.

Future<int> sumStream(Stream<int> stream) async {
  var sum = 0;
  await for (final value in stream) {
    sum += value;
  }
  return sum;
}

Stream<int> countStream(int to) async* {
  for (int i = 1; i <= to; i++) {
    yield i;
  }
}

void main() async {
  var stream = countStream(10);
  var sum = await sumStream(stream);
  print(sum); // 55
}

에러 이벤트

스트림에 더 이상 이벤트가 없으면 스트림은 끝나는데, 이벤트를 받는 코드는 새 이벤트가 도착할 때처럼 이것도 통지를 받아요. await for 루프로 이벤트를 읽을 때는 스트림이 끝나면 루프도 멈춥니다.

스트림은 데이터 이벤트처럼 에러 이벤트도 전달할 수 있어요. 대부분의 스트림은 첫 에러 이후에 멈추지만, 에러를 여러 번 전달하거나 에러 이벤트 뒤에도 데이터를 더 전달하는 스트림도 가능해요. 이 문서에서는 최대 하나의 에러만 전달하는 스트림만 다룰게요.

await for로 스트림을 읽을 때는 루프 문이 에러를 던지면서 루프도 함께 끝나요. try-catch로 그 에러를 잡을 수 있어요. 다음 예시는 루프 반복자가 4일 때 에러를 던집니다.

Future<int> sumStream(Stream<int> stream) async {
  var sum = 0;
  try {
    await for (final value in stream) {
      sum += value;
    }
  } catch (e) {
    return -1;
  }
  return sum;
}

Stream<int> countStream(int to) async* {
  for (int i = 1; i <= to; i++) {
    if (i == 4) {
      throw Exception('Intentional exception');
    } else {
      yield i;
    }
  }
}

void main() async {
  var stream = countStream(10);
  var sum = await sumStream(stream);
  print(sum); // -1
}

스트림 다루기

Stream 클래스에는 Iterable의 메서드와 비슷하게 스트림에 공통 작업을 해 주는 도우미 메서드가 많이 있어요. 예를 들어 스트림에서 마지막 양의 정수를 찾으려면 Stream API의 lastWhere()를 씁니다.

Future<int> lastPositive(Stream<int> stream) => stream.lastWhere((x) => x >= 0);

두 종류의 스트림

스트림에는 두 종류가 있어요.

단일 구독 스트림

가장 흔한 스트림은 더 큰 전체의 일부인 이벤트들의 연속이에요. 이벤트는 순서대로, 하나도 빠짐없이 전달되어야 합니다. 파일을 읽거나 웹 요청을 받을 때 얻는 스트림이 바로 이 종류예요. 이런 스트림은 한 번만 들을 수 있어요. 다시 듣기 시작하면 처음 이벤트를 놓쳐서 나머지 스트림이 의미 없어질 수 있기 때문이에요. 듣기를 시작하면 데이터를 청크 단위로 가져와 제공합니다.

브로드캐스트 스트림

다른 종류의 스트림은 하나씩 따로 처리할 수 있는 개별 메시지를 위한 것이에요. 예를 들어 브라우저의 마우스 이벤트에 쓸 수 있어요. 이런 스트림은 아무 때나 듣기 시작할 수 있고, 듣는 동안 발생한 이벤트를 받습니다. 여러 리스너가 동시에 들을 수 있고, 이전 구독을 취소한 뒤 나중에 다시 들을 수도 있어요.

스트림을 처리하는 메서드

Stream<T>에서 스트림을 처리해 결과를 돌려주는 메서드들은 다음과 같아요.

Future<T> get first;
Future<bool> get isEmpty;
Future<T> get last;
Future<int> get length;
Future<T> get single;
Future<bool> any(bool Function(T element) test);
Future<bool> contains(Object? needle);
Future<E> drain<E>([E? futureValue]);
Future<T> elementAt(int index);
Future<bool> every(bool Function(T element) test);
Future<T> firstWhere(bool Function(T element) test, {T Function()? orElse});
Future<S> fold<S>(S initialValue, S Function(S previous, T element) combine);
Future forEach(void Function(T element) action);
Future<String> join([String separator = ', ']);
Future<T> lastWhere(bool Function(T element) test, {T Function()? orElse});
Future pipe(StreamConsumer<T> streamConsumer);
Future<T> reduce(T Function(T previous, T element) combine);
Future<T> singleWhere(bool Function(T element) test, {T Function()? orElse});
Future<List<T>> toList();
Future<Set<T>> toSet();

이 메서드들은 drain()pipe()를 빼고 모두 Iterable의 비슷한 함수에 대응해요. 각각은 async 함수와 await for 루프(또는 다른 메서드 하나)로 쉽게 작성할 수 있어요.

스트림을 수정하는 메서드

Stream에서 원래 스트림을 바탕으로 새 스트림을 반환하는 메서드도 있어요. 각각은 새 스트림에 누군가 듣기 시작할 때까지 원래 스트림 듣기를 기다립니다.

Stream<R> cast<R>();
Stream<S> expand<S>(Iterable<S> Function(T element) convert);
Stream<S> map<S>(S Function(T event) convert);
Stream<T> skip(int count);
Stream<T> skipWhile(bool Function(T element) test);
Stream<T> take(int count);
Stream<T> takeWhile(bool Function(T element) test);
Stream<T> where(bool Function(T event) test);

이 메서드들은 iterable을 다른 iterable로 바꾸는 Iterable의 비슷한 메서드에 대응해요. 모두 async 함수와 await for 루프로 쉽게 작성할 수 있습니다.

Stream<E> asyncExpand<E>(Stream<E>? Function(T event) convert);
Stream<E> asyncMap<E>(FutureOr<E> Function(T event) convert);
Stream<T> distinct([bool Function(T previous, T next)? equals]);

asyncExpand()asyncMap()expand()map()과 비슷하지만 함수 인자를 비동기 함수로 받을 수 있어요. distinct()는 Iterable에는 없지만 있을 법한 함수예요.

Stream<T> handleError(Function onError, {bool Function(dynamic error)? test});
Stream<T> timeout(Duration timeLimit, {void Function(EventSink<T> sink)? onTimeout});
Stream<S> transform<S>(StreamTransformer<T, S> streamTransformer);

마지막 세 함수는 더 특화되어 있어요. 이들은 await for 루프로 직접 관리할 수 없는 에러 처리를 다룹니다. 루프는 첫 에러를 만나면 루프와 스트림 구독을 함께 끝내고, 복구할 내장 메커니즘이 없어요. 다음 코드는 스트림이 소비되기 전에 handleError()로 에러를 걸러 내는 방법을 보여줍니다.

Stream<S> mapLogErrors<S, T>(
  Stream<T> stream,
  S Function(T event) convert,
) async* {
  var streamWithoutErrors = stream.handleError((e) => log(e));
  var streamWithTimeout = streamWithoutErrors.timeout(
    const Duration(seconds: 2),
    onTimeout: (eventSink) {
      eventSink.addError('Timed out after 2 seconds');
      eventSink.close();
    },
  );

  await for (final event in streamWithTimeout) {
    yield convert(event);
  }
}

transform() 함수

transform() 함수는 에러 처리만을 위한 게 아니에요. 스트림에 대한 더 일반화된 "map"이라고 볼 수 있어요. 보통의 map은 들어오는 이벤트마다 값 하나를 요구하지만, 특히 I/O 스트림에서는 출력 이벤트 하나를 만들기 위해 여러 입력 이벤트가 필요할 수 있어요. StreamTransformer가 바로 그런 작업을 처리할 수 있어요. 예를 들어 Utf8Decoder 같은 디코더는 트랜스포머예요. 트랜스포머는 bind() 함수 하나만 구현하면 되는데, 이 함수는 async 함수로 쉽게 구현할 수 있어요.

파일 읽고 디코딩하기

다음 코드는 파일을 읽고 스트림에 두 개의 트랜스폼을 적용해요. 먼저 데이터를 UTF8로 변환하고, 그다음 LineSplitter를 통과시킵니다. #으로 시작하는 줄을 제외하고 모든 줄을 출력해요.

import 'dart:convert';
import 'dart:io';

void main(List<String> args) async {
  var file = File(args[0]);
  var lines = utf8.decoder
      .bind(file.openRead())
      .transform(const LineSplitter());
  await for (final line in lines) {
    if (!line.startsWith('#')) {
      print(line);
    }
  }
}

listen() 메서드

Stream의 마지막 메서드는 listen()이에요. 이것은 "저수준" 메서드로, 다른 모든 스트림 함수는 listen()을 기준으로 정의됩니다.

StreamSubscription<T> listen(void Function(T event)? onData, {Function? onError, void Function()? onDone, bool? cancelOnError});

Stream 타입을 만들려면 Stream 클래스를 확장하고 listen() 메서드만 구현하면 돼요. Stream의 다른 메서드는 모두 동작을 위해 listen()을 호출합니다.

listen() 메서드는 스트림 듣기를 시작하게 해요. 듣기 전까지 스트림은 어떤 이벤트를 보려는지 설명하는 비활성 객체일 뿐이에요. 들을 때는 이벤트를 만들어내는 활성 스트림을 나타내는 StreamSubscription 객체가 반환됩니다. 이것은 Iterable이 그저 객체 모음일 뿐이고 실제 순회는 iterator가 하는 것과 비슷한 관계예요. 스트림 구독으로 구독을 일시 정지하고, 재개하고, 완전히 취소할 수 있어요. 데이터 이벤트나 에러 이벤트마다, 그리고 스트림이 닫힐 때 호출할 콜백을 지정할 수 있어요.

더 알아보기