Spec-Zone.ru › Dart 2

dart:async

Класс Stream<T>

Источник асинхронных событий данных.

Поток предоставляет способ получения последовательности событий. Каждое событие является либо событием данных, также называемым элементом потока, либо событием об ошибке, которое является уведомлением о том, что что-то не удалось. Когда поток исчерпал все свои события, одно событие «завершения» уведомляет слушателя о достижении конца.

Вы создаете поток, вызвав функцию async*, которая затем возвращает поток. Потребление этого потока заставит функцию генерировать события до завершения, и поток закроется. Вы потребляете поток либо с помощью цикла await for, который доступен внутри функции async или async*, либо, перенаправляя его события напрямую, используя yield* внутри функции async*. Пример:

Stream<T> optionalMap<T>(
    Stream<T> source , [T Function(T)? convert]) async* {
  if (convert == null) {
    yield* source;
  } else {
    await for (var event in source) {
      yield convert(event);
    }
  }
}

Когда эта функция вызывается, она сразу же возвращает объект Stream<T>. Затем ничего не происходит, пока кто-то не попытается потреблять этот поток. В этот момент начинает выполняться тело функции async*. Если функция convert была опущена, то yield* будет прослушивать поток source и перенаправлять все события, данные и ошибки, в возвращаемый поток. Когда поток source завершается, выполнение yield* завершается, и тело функции optionalMap тоже завершается. Это закрывает возвращаемый поток. Если convert есть, функция вместо этого прослушивает исходный поток и входит в цикл await for, который многократно ожидает следующего события данных. При получении события данных она вызывает convert со значением и отправляет результат в возвращаемый поток. Если поток source не отправляет событий об ошибках, цикл завершается, когда поток source завершается, а затем завершается тело функции optionalMap, что закрывает возвращаемый поток. При получении события об ошибке от потока source функция await for повторно выбрасывает эту ошибку, что прерывает цикл. Затем ошибка достигает конца тела функции optionalMap, поскольку она не перехвачена. Это приводит к тому, что ошибка отправляется в возвращаемый поток, который затем закрывается.

Класс Stream также предоставляет функциональность, которая позволяет вам вручную прослушивать события из потока или преобразовывать поток в другой поток или в будущее.

Функция forEach соответствует циклу await for, так же как Iterable.forEach соответствует обычному циклу for/in. Как и в цикле, она будет вызывать функцию для каждого события данных и прерываться на ошибке.

Более низкоуровневый метод listen является основой для всех других методов. Вы вызываете listen на потоке, чтобы указать, что хотите получать события, и зарегистрировать обратные вызовы, которые будут получать эти события. Когда вы вызываете listen, вы получаете объект StreamSubscription, который является активным объектом, предоставляющим события, и который может быть использован для остановки прослушивания или временной приостановки событий от подписки.

Существует два типа потоков: потоки «с единственной подпиской» и «вещательные» потоки.

Поток с единственной подпиской позволяет только одному слушателю в течение всего жизненного цикла потока. Он не начинает генерировать события до тех пор, пока не будет слушателя, и прекращает отправку событий, когда слушатель отписывается, даже если источник событий по-прежнему может предоставить больше. Поток, созданный функцией async*, является потоком с единственной подпиской, но каждый вызов функции создает новый такой поток.

Двойное прослушивание потока с единственной подпиской не допускается, даже после отмены первой подписки.

Потоки с единственной подпиской обычно используются для потоковой передачи фрагментов больших непрерывных данных, например, при работе с файловым вводом-выводом.

Вещательный поток позволяет любое количество слушателей, и он отправляет свои события, когда они готовы, независимо от того, есть ли слушатели или нет.

Вещательные потоки используются для независимых событий/наблюдателей.

Если несколько слушателей хотят прослушивать поток с единственной подпиской, используйте asBroadcastStream для создания вещательного потока на основе не-вещательного потока.

Для любого типа потока преобразования потока, такие как where и skip, возвращают тот же тип потока, что и тот, на котором был вызван метод, если не указано иное.

Когда генерируется событие, слушатели в этот момент получат событие. Если слушатель добавляется к вещательному потоку во время отправки события, этот слушатель не получит событие, которое сейчас отправляется. Если слушатель отменяется, он немедленно прекращает получение событий. Прослушивание вещательного потока можно рассматривать как прослушивание нового потока, содержащего только события, которые еще не были отправлены во время вызова listen. Например, свойство first прослушивает поток, а затем возвращает первое событие, которое получает слушатель. Это не обязательно первое событие, отправленное потоком, но первое из оставшихся событий вещательного потока.

Когда отправляется событие «завершения», подписчики отписываются перед получением события. После отправки события поток не имеет подписчиков. Добавление новых подписчиков к вещательному потоку после этого момента разрешено, но они просто получат новое событие «завершения» как можно скорее.

Подписки на поток всегда учитывают запросы «паузы». При необходимости они должны буферизовать свой ввод, но часто, и предпочтительно, они могут просто запросить паузу своего ввода тоже.

По умолчанию реализация isBroadcast возвращает false. Вещательный поток, наследующий от Stream, должен переопределять isBroadcast, чтобы вернуть true, если он хочет указать, что он ведет себя как вещательный поток.

Реализации
  • CustomStream
  • ElementStream
  • HttpClientResponse
  • HttpRequest
  • HttpServer
  • RawDatagramSocket
  • RawSecureServerSocket
  • RawServerSocket
  • RawSocket
  • ReceivePort
  • SecureServerSocket
  • ServerSocket
  • Socket
  • Stdin
  • StreamView
  • WebSocket

Конструкторы

Stream()
Stream.empty()
const
factory
Создает пустой вещательный поток.
Stream.error(Object error, [StackTrace? stackTrace])
factory
Создает поток, который отправляет одно событие об ошибке перед завершением.
Stream.eventTransformed(Stream source, EventSink mapSink(EventSink<T> sink))
factory
Создает поток, где все события существующего потока проходят через преобразование источника.
Stream.fromFuture(Future<T> future)
factory
Создает новый поток с единственной подпиской из будущего.
Stream.fromFutures(Iterable<Future<T>> futures)
factory
Создает поток с единственной подпиской из группы будущих значений.
Stream.fromIterable(Iterable<T> elements)
factory
Создает поток, который получает данные из elements.
Stream.multi(void onListen(MultiStreamController<T>), {bool isBroadcast = false})
factory
Создает поток с множественной подпиской.
Stream.periodic(Duration period, [T computation(int computationCount)?])
factory
Создает поток, который повторяет отправку событий через period интервалы.
Stream.value(T value)
factory
Создает поток, который отправляет одно событие данных перед закрытием.

Свойства

первый → Future<T>
только для чтения
Первый элемент этого потока.
hashCode → int
только для чтения, унаследованный
Код хэша для этого объекта.
isBroadcast → bool
только для чтения
Является ли этот поток потоком широковещательной передачи.
isEmpty → Future<bool>
только для чтения
Содержит ли этот поток какие-либо элементы.
последний → Future<T>
только для чтения
Последний элемент этого потока.
длина → Future<int>
только для чтения
Количество элементов в этом потоке.
runtimeType → Type
только для чтения, унаследованный
Представление типа объекта во время выполнения.
единственный → Future<T>
только для чтения
Единственный элемент этого потока.

Методы

any(bool test(T element)) → Future<bool>
Проверяет, принимает ли test любой элемент, предоставленный этим потоком.
asBroadcastStream({void onListen(StreamSubscription<T> subscription)?, void onCancel(StreamSubscription<T> subscription)?}) → Stream<T>
Возвращает поток с множественным подписанием, который генерирует те же события, что и этот.
asyncExpand<E>(Stream<E>? convert(T event)) → Stream<E>
Преобразует каждый элемент в последовательность асинхронных событий.
asyncMap<E>(FutureOr<E> convert(T event)) → Stream<E>
Создает новый поток, где каждое событие данных этого потока асинхронно отображается на новое событие.
cast<R>() → Stream<R>
Адаптирует этот поток для Stream<R>.
contains(Object? needle) → Future<bool>
Возвращает, присутствует ли needle в элементах, предоставляемых этим потоком.
distinct([bool equals(T previous, T next)?]) → Stream<T>
Пропускает события данных, если они равны предыдущему событию данных.
drain<E>([E? futureValue]) → Future<E>
Отбрасывает все данные из этого потока, но сигнализирует о завершении или возникновении ошибки.
elementAt(int index) → Future<T>
Возвращает значение index-го события данных этого потока.
every(bool test(T element)) → Future<bool>
Проверяет, принимает ли test все элементы, предоставленные этим потоком.
expand<S>(Iterable<S> convert(T element)) → Stream<S>
Преобразует каждый элемент этого потока в последовательность элементов.
firstWhere(bool test(T element), {T orElse()?}) → Future<T>
Находит первый элемент этого потока, соответствующий test.
fold<S>(S initialValue, S combine(S previous, T element)) → Future<S>
Объединяет последовательность значений, многократно применяя combine.
forEach(void action(T element)) → Future
Выполняет action для каждого элемента этого потока.
handleError(Function onError, {bool test(dynamic error)?}) → Stream<T>
Создаёт обернутый поток, который перехватывает некоторые ошибки из этого потока.
join([String separator = ""]) → Future<String>
Объединяет строковое представление элементов в одну строку.
lastWhere(bool test(T element), {T orElse()?}) → Future<T>
Находит последний элемент в этом потоке, соответствующий test.
listen(void onData(T event)?, {Function? onError, void onDone()?, bool? cancelOnError}) → StreamSubscription<T>
Добавляет подписку на этот поток.
map<S>(S convert(T event)) → Stream<S>
Преобразует каждый элемент этого потока в новое событие потока.
noSuchMethod(Invocation invocation) → dynamic
унаследованный
Вызывается при обращении к несуществующему методу или свойству.
pipe(StreamConsumer<T> streamConsumer) → Future
Перенаправляет события этого потока в streamConsumer.
reduce(T combine(T previous, T element)) → Future<T>
Объединяет последовательность значений, многократно применяя combine.
singleWhere(bool test(T element), {T orElse()?}) → Future<T>
Находит единственный элемент в этом потоке, соответствующий test.
skip(int count) → Stream<T>
Пропускает первые count событий данных из этого потока.
skipWhile(bool test(T element)) → Stream<T>
Пропускает события данных из этого потока, пока они соответствуют test.
take(int count) → Stream<T>
Предоставляет не более первых count событий данных этого потока.
takeWhile(bool test(T element)) → Stream<T>
Передает события данных, пока test успешно.
timeout(Duration timeLimit, {void onTimeout(EventSink<T> sink)?}) → Stream<T>
Создаёт новый поток с теми же событиями, что и этот поток.
toList() → Future<List<T>>
Собрать все элементы этого потока в List.
toSet() → Future<Set<T>>
Собрать данные этого потока в Set.
toString() → String
унаследованный
Строковое представление этого объекта.
transform<S>(StreamTransformer<T, S> streamTransformer) → Stream<S>
Применяет streamTransformer к этому потоку.
where(bool test(T event)) → Stream<T>
Создает новый поток из этого потока, отбрасывая некоторые элементы.

Операторы

operator ==(Object other) → bool
унаследованный
Оператор равенства.

Статические методы

castFrom<S, T>(Stream<S> source) → Stream<T>
Адаптирует source для Stream<T>.

© 2012 the Dart project authors
Licensed under the BSD 3-Clause "New" or "Revised" License.
https://api.dart.dev/stable/2.18.5/dart-async/Stream-class.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API