Класс 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, если он хочет указать, что он ведет себя как вещательный поток.
- Реализации
Конструкторы
- Stream()
- Stream.empty() constfactory
- Создает пустой вещательный поток.
- 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