Spec-Zone.ru › ReactiveX

Buffer

периодически собирает элементы, выпущенные Observable, в пакеты и выдает эти пакеты, а не выдает элементы по одному
Buffer

Оператор Buffer преобразует Observable, который выдает элементы, в Observable, который выдает буферизованные коллекции этих элементов. В различных реализациях Buffer, специфичных для языка, существует ряд вариантов, которые различаются по способу выбора элементов для включения в буферы.

Обратите внимание, что если исходный Observable выдает уведомление onError, Buffer немедленно передаст это уведомление, не выпустив предварительно буфер, который он формирует, даже если этот буфер содержит элементы, выпущенные исходным Observable до выдачи им уведомления об ошибке.

Оператор Window похож на Buffer, но собирает элементы в отдельные Observables, а не в структуры данных, прежде чем повторно выпустить их.

См. также

  • Window
  • Введение в Rx: Buffer
  • Введение в Rx: Buffer revisited
  • 101 пример Rx: Buffer — Simple

Информация, специфичная для языка

RxCpp buffer pairwise

RxCpp реализует два варианта Buffer:

buffer(count)

buffer(count)

buffer(count) выдает неперекрывающиеся буферы в виде vectors, каждый из которых содержит не более count элементов из исходного Observable (конечный выпущенный vector может содержать меньше, чем count элементов).

buffer(count, skip)

buffer(count,skip)

buffer(count, skip) создает новый буфер, начиная с первого выпущенного элемента из исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими. Он выдает эти буферы как vectors. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент) или иметь пробелы (где элементы, выпущенные исходным Observable, не представлены ни в одном буфере).

RxGroovy buffer

В RxGroovy есть несколько вариантов Buffer:

buffer(count)

buffer(count)

buffer(count) выдает неперекрывающиеся буферы в виде Lists, каждый из которых содержит не более count элементов из исходного Observable (конечный выпущенный List может содержать меньше, чем count элементов).

  • Javadoc: buffer(int)

buffer(count, skip)

buffer(count,skip)

buffer(count, skip) создает новый буфер, начиная с первого выпущенного элемента из исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими. Он выдает эти буферы как Lists. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент) или иметь пробелы (где элементы, выпущенные исходным Observable, не представлены ни в одном буфере).

  • Javadoc: buffer(int,int)

buffer(bufferClosingSelector)

buffer(bufferClosingSelector)

Когда он подписывается на исходный Observable, buffer(bufferClosingSelector) начинает собирать его излучения в List, а также вызывает bufferClosingSelector для генерации второго Observable. Когда этот второй Observable выдает объект TClosing, buffer выдает текущий List и повторяет этот процесс: начинает новый List и вызывает bufferClosingSelector для создания нового Observable для мониторинга. Он будет делать это до тех пор, пока исходный Observable не завершится.

  • Javadoc: buffer(Func0)

buffer(boundary[, initialCapacity])

buffer(boundary)

buffer(boundary) отслеживает Observable, boundary. Каждый раз, когда этот Observable выдает элемент, он создает новый List для начала сбора элементов, выпущенных исходным Observable, и выдает предыдущий List.

  • Javadoc: buffer(Observable)
  • Javadoc: buffer(Observable,int)

buffer(bufferOpenings, bufferClosingSelector)

buffer(bufferOpenings,bufferClosingSelector)

buffer(bufferOpenings, bufferClosingSelector) отслеживает Observable, bufferOpenings, который выдает объекты BufferOpening. Каждый раз, когда он наблюдает такой выпущенный элемент, он создает новый List для начала сбора элементов, выпущенных исходным Observable, и передает Observable bufferOpenings в функцию closingSelector. Эта функция возвращает Observable. buffer отслеживает этот Observable, и когда он обнаруживает выпущенный элемент от него, он закрывает List и выдает его как свое собственное излучение.

  • Javadoc: buffer(Observable,Func1)

buffer(timespan, unit[, scheduler])

buffer(timespan,unit)

buffer(timespan, unit) выдает новый List элементов периодически, каждые timespan единиц времени, содержащий все элементы, выпущенные исходным Observable с момента предыдущего выпуска пакета или, в случае первого пакета, с момента подписки на исходный Observable. Существует также версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,TimeUnit)
  • Javadoc: buffer(long,TimeUnit,Scheduler)

buffer(timespan, unit, count[, scheduler])

buffer(timespan,unit,count)

buffer(timespan, unit, count) выдает новый List элементов для каждых count элементов, выпущенных исходным Observable, или, если прошло timespan с момента его последнего выпуска пакета, он выдает пакет из любого количества элементов, которые исходный Observable выпустил за этот период, даже если это меньше, чем count. Существует также версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,TimeUnit,int)
  • Javadoc: buffer(long,TimeUnit,int,Scheduler)

buffer(timespan, timeshift, unit[, scheduler])

buffer(timespan,timeshift,unit)

buffer(timespan, timeshift, unit) создает новый List элементов каждые timeshift период времени и заполняет этот пакет каждым элементом, выпущенным исходным Observable с этого момента до тех пор, пока не пройдет timespan времени с момента создания пакета, прежде чем выпустить этот List как свое собственное излучение. Если timespan длиннее, чем timeshift, выпущенные пакеты будут представлять перекрывающиеся периоды времени, поэтому они могут содержать повторяющиеся элементы. Существует также версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,long,TimeUnit)
  • Javadoc: buffer(long,long,TimeUnit,Scheduler)

Вы можете использовать оператор Buffer для реализации противодавления (то есть для работы с Observable, который может производить элементы слишком быстро для их потребления наблюдателем).

Buffer as a backpressure strategy

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

Пример кода

Observable<List<Integer>> burstyBuffered = bursty.buffer(500, TimeUnit.MILLISECONDS);
Buffer as a backpressure strategy

Или вы можете сделать это сложнее и собирать элементы в буферы во время импульсных периодов и выдавать их в конце каждого импульса, используя оператор Debounce для выдачи индикатора закрытия буфера оператору буфера.

Пример кода

// we have to multicast the original bursty Observable so we can use it
// both as our source and as the source for our buffer closing selector:
Observable<Integer> burstyMulticast = bursty.publish().refCount();
// burstyDebounced will be our buffer closing selector:
Observable<Integer> burstyDebounced = burstyMulticast.debounce(10, TimeUnit.MILLISECONDS);
// and this, finally, is the Observable of buffers we're interested in:
Observable<List<Integer>> burstyBuffered = burstyMulticast.buffer(burstyDebounced);

RxJava 1․x buffer

В RxJava существуют несколько вариантов оператора Buffer:

buffer(count)

buffer(count)

buffer(count) излучает неперекрывающиеся буферы в виде Listы, каждый из которых содержит не более count элементов из исходного Observable (окончательный излученный List может содержать меньше, чем count элементов).

  • Javadoc: buffer(int)

buffer(count, skip)

buffer(count,skip)

buffer(count, skip) создает новый буфер, начиная с первого излученного элемента из исходного Observable, и каждый skip элемент после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими. Он излучает эти буферы в виде Listы. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент) или иметь разрывы (где элементы, излученные исходным Observable, не представлены в любом буфере).

  • Javadoc: buffer(int,int)

buffer(bufferClosingSelector)

buffer(bufferClosingSelector)

При подписке на исходный Observable, buffer(bufferClosingSelector) начинает собирать его излучения в List, и также вызывает bufferClosingSelector для генерации второго Observable. Когда этот второй Observable излучает объект TClosing, buffer излучает текущий List и повторяет этот процесс: начинает новый List и вызывает bufferClosingSelector для создания нового Observable для мониторинга. Он будет делать это до тех пор, пока исходный Observable не завершится.

  • Javadoc: buffer(Func0)

buffer(boundary)

buffer(boundary)

buffer(boundary) отслеживает Observable, boundary. Каждый раз, когда этот Observable излучает элемент, он создает новый List для начала сбора элементов, излучаемых исходным Observable, и излучает предыдущий List.

  • Javadoc: buffer(Observable)
  • Javadoc: buffer(Observable,int)

buffer(bufferOpenings, bufferClosingSelector)

buffer(bufferOpenings,bufferClosingSelector)

buffer(bufferOpenings, bufferClosingSelector) отслеживает Observable, bufferOpenings, который излучает BufferOpening объекты. Каждый раз, когда он наблюдает такой излученный элемент, он создает новый List для начала сбора элементов, излучаемых исходным Observable, и передает Observable bufferOpenings в функцию closingSelector. Эта функция возвращает Observable. buffer отслеживает этот Observable, и когда он обнаруживает излученный элемент от него, он закрывает List и излучает его как свое собственное излучение.

  • Javadoc: buffer(Observable,Func1)

buffer(timespan, unit[, scheduler])

buffer(timespan,unit)

buffer(timespan, unit) излучает новый List элементов периодически, каждые timespan единиц времени, содержащий все элементы, излученные исходным Observable с момента предыдущего излучения пакета или, в случае первого пакета, с момента подписки на исходный Observable. Также существует вариант этого варианта оператора, который принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,TimeUnit)
  • Javadoc: buffer(long,TimeUnit,Scheduler)

buffer(timespan, unit, count[, scheduler])

buffer(timespan,unit,count)

buffer(timespan, unit, count) излучает новый List элементов для каждого count элементов, излученных исходным Observable, или, если timespan прошло с момента последнего излучения пакета, он излучает пакет из любого количества элементов, которые излучил исходный Observable за этот период, даже если это меньше count. Также существует вариант этого варианта оператора, который принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,TimeUnit,int)
  • Javadoc: buffer(long,TimeUnit,int,Scheduler)

buffer(timespan, timeshift, unit[, scheduler])

buffer(timespan,timeshift,unit)

buffer(timespan, timeshift, unit) создает новый List элементов каждые timeshift единицы времени и заполняет этот пакет каждым элементом, излученным исходным Observable с этого момента до тех пор, пока не пройдет timespan единиц времени с момента создания пакета, прежде чем излучить этот List в качестве собственного излучения. Если timespan длиннее, чем timeshift, излучаемые пакеты будут представлять перекрывающиеся временные периоды, поэтому они могут содержать дублируемые элементы. Также существует вариант этого варианта оператора, который принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

  • Javadoc: buffer(long,long,TimeUnit)
  • Javadoc: buffer(long,long,TimeUnit,Scheduler)

Вы можете использовать оператор Buffer для реализации обратной связи (то есть, для работы с Observable, который может генерировать элементы слишком быстро для того, чтобы его наблюдатель их потреблял).

Buffer as a backpressure strategy

Buffer может уменьшить последовательность многих элементов до последовательности меньшего количества буферов элементов, делая их более управляемыми. Например, вы могли бы периодически закрывать и излучать буфер элементов из бурного Observable с регулярным временным интервалом.

Пример кода

Observable<List<Integer>> burstyBuffered = bursty.buffer(500, TimeUnit.MILLISECONDS);
Buffer as a backpressure strategy

Или вы могли бы усложнить ситуацию, и собирать элементы в буферах во время бурных периодов и излучать их в конце каждого выброса, используя оператор Debounce для излучения индикатора закрытия буфера в оператор буферизации.

Пример кода

// we have to multicast the original bursty Observable so we can use it
// both as our source and as the source for our buffer closing selector:
Observable<Integer> burstyMulticast = bursty.publish().refCount();
// burstyDebounced will be our buffer closing selector:
Observable<Integer> burstyDebounced = burstyMulticast.debounce(10, TimeUnit.MILLISECONDS);
// and this, finally, is the Observable of buffers we're interested in:
Observable<List<Integer>> burstyBuffered = burstyMulticast.buffer(burstyDebounced);

См. также

  • DebouncedBuffer With RxJava by Gopal Kaushik
  • DebounceBuffer: Use publish(), debounce() and buffer() together to capture bursts of events. by Ben Christensen

RxJS buffer bufferWithCount bufferWithTime bufferWithTimeOrCount

RxJS имеет четыре оператора Buffer — buffer, bufferWithCount, bufferWithTime, и bufferWithTimeOrCount — каждый из которых имеет варианты, определяющие различные способы управления тем, какие элементы источника Observable будут испускаться в рамках каких буферов.

buffer(bufferBoundaries)

buffer(bufferBoundaries)

buffer(bufferBoundaries) отслеживает Observable, bufferBoundaries. Каждый раз, когда этот Observable испускает элемент, он создает новую коллекцию для начала сбора элементов, испускаемых исходным Observable, и испускает предыдущую коллекцию.

buffer(bufferClosingSelector)

buffer(bufferClosingSelector)

При подписке на исходный Observable, buffer(bufferClosingSelector) начинает собирать его испускания в коллекцию, а также вызывает bufferClosingSelector для генерации второго Observable. Когда этот второй Observable испускает элемент, buffer испускает текущую коллекцию и повторяет этот процесс: начинает новую коллекцию и вызывает bufferClosingSelector для создания нового Observable для отслеживания. Он будет делать это до тех пор, пока исходный Observable не завершит работу.

buffer(bufferOpenings,bufferClosingSelector)

buffer(bufferOpenings,bufferClosingSelector)

buffer(bufferOpenings, bufferClosingSelector) отслеживает Observable, bufferOpenings, который испускает BufferOpening объекты. Каждый раз, когда он наблюдает такой испущенный элемент, он создает новую коллекцию для начала сбора элементов, испускаемых исходным Observable, и передаёт Observable bufferOpenings в функцию bufferClosingSelector. Эта функция возвращает Observable. buffer отслеживает этот Observable, и когда обнаруживает испущенный элемент, испускает текущую коллекцию и начинает новую.

buffer присутствует в следующих дистрибутивах:

  • rx.all.js
  • rx.all.compat.js
  • rx.coincidence.js

buffer требует одного из следующих дистрибутивов:

  • rx.js
  • rx.compat.js
  • rx.lite.js
  • rx.lite.compat.js

bufferWithCount(count)

bufferWithCount(count)

bufferWithCount(count) испускает неперекрывающиеся буферы, каждый из которых содержит не более count элементов из исходного Observable (последний испущенный буфер может содержать меньше, чем count элементов).

bufferWithCount(count, skip)

bufferWithCount(count,skip)

bufferWithCount(count, skip) создаёт новый буфер, начиная с первого испущенного элемента исходного Observable, и новый для каждого skip элемента после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими, испуская каждый буфер при его завершении. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент) или иметь разрывы (где элементы, испущенные исходным Observable, не представлены ни в одном буфере).

bufferWithCount присутствует в следующих дистрибутивах:

  • rx.js
  • rx.compat.js
  • rx.all.js
  • rx.all.compat.js
  • rx.lite.extras.js

bufferWithTime(timeSpan)

bufferWithTime(timeSpan)

bufferWithTime(timeSpan) испускает новую коллекцию элементов периодически, каждые timeSpan миллисекунд, содержащую все элементы, испущенные исходным Observable с момента предыдущего испускания пакета или, в случае первого пакета, с момента подписки на исходный Observable. Также есть версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию используется планировщик timeout.

bufferWithTime(timeSpan, timeShift)

bufferWithTime(timeSpan,timeShift)

bufferWithTime(timeSpan, timeShift) создает новую коллекцию элементов каждые timeShift миллисекунд и заполняет этот пакет каждым элементом, испущенным исходным Observable с этого момента до timeSpan миллисекунд после создания коллекции, прежде чем испустить эту коллекцию как собственное испускание. Если timeSpan больше, чем timeShift, испущенные пакеты будут представлять перекрывающиеся временные интервалы, и, следовательно, они могут содержать дубликаты элементов. Также есть версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию используется планировщик timeout.

bufferWithTimeOrCount(timeSpan, count)

bufferWithTimeOrCount(timeSpan,count)

bufferWithTimeOrCount(timeSpan, count) испускает новую коллекцию элементов для каждых count элементов, испущенных исходным Observable, или, если прошло timeSpan миллисекунд с момента последнего испускания коллекции, он испускает коллекцию из того количества элементов, которое исходный Observable испустил за этот промежуток времени, даже если это меньше count. Также есть версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию используется планировщик timeout.

bufferWithTime и bufferWithTimeOrCount присутствуют в следующих дистрибутивах:

  • rx.all.js
  • rx.all.compat.js
  • rx.time.js

bufferWithTime и bufferWithTimeOrCount требуют одного из следующих дистрибутивов:

  • rx.time.js требует rx.js или rx.compat.js
  • в противном случае: rx.lite.js или rx.lite.compat.js

RxKotlin buffer

В RxKotlin существует несколько вариантов Buffer:

buffer(count)

buffer(count)

buffer(count) генерирует непересекающиеся буферы в виде List, каждый из которых содержит не более count элементов из исходного Observable (последний выпущенный List может содержать меньше, чем count элементов).

buffer(count, skip)

buffer(count,skip)

buffer(count, skip) создаёт новый буфер, начиная с первого элемента, выпущенного из исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими. Эти буферы излучаются как List. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент), или у них могут быть пробелы (где элементы, выпущенные исходным Observable, не представлены ни в одном буфере).

buffer(bufferClosingSelector)

buffer(bufferClosingSelector)

При подписке на исходный Observable, buffer(bufferClosingSelector) начинает собирать его излучения в List, а также вызывает bufferClosingSelector для генерации второго Observable. Когда этот второй Observable излучает объект TClosing, buffer излучает текущий List и повторяет этот процесс: начинает новый List и вызывает bufferClosingSelector для создания нового Observable для мониторинга. Это будет продолжаться до завершения исходного Observable.

buffer(boundary)

buffer(boundary)

buffer(boundary) отслеживает Observable, boundary. Каждый раз, когда этот Observable излучает элемент, он создаёт новый List для начала сбора элементов, излучаемых исходным Observable, и излучает предыдущий List.

buffer(bufferOpenings, bufferClosingSelector)

buffer(bufferOpenings,bufferClosingSelector)

buffer(bufferOpenings, bufferClosingSelector) отслеживает Observable, bufferOpenings, который излучает BufferOpening объекты. Каждый раз, когда он наблюдает такой выпущенный элемент, он создает новый List для начала сбора элементов, излучаемых исходным Observable, и передает bufferOpenings Observable в функцию closingSelector. Эта функция возвращает Observable. buffer отслеживает этот Observable, и когда он обнаруживает выпущенный элемент из него, он закрывает List и излучает его как своё собственное излучение.

buffer(timespan, unit[, scheduler])

buffer(timespan,unit)

buffer(timespan, unit) излучает новый List элементов периодически, каждые timespan единиц времени, содержащий все элементы, излученные исходным Observable с момента предыдущего излучения пакета или, в случае первого пакета, с момента подписки на исходный Observable. Также существует версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

buffer(timespan, unit, count[, scheduler])

buffer(timespan,unit,count)

buffer(timespan, unit, count) излучает новый List элементов для каждого count элементов, излученных исходным Observable, или, если timespan истек с момента последнего излучения пакета, он излучает пакет из количества элементов, которые исходный Observable излучил за этот промежуток времени, даже если это меньше count. Также существует версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

buffer(timespan, timeshift, unit[, scheduler])

buffer(timespan,timeshift,unit)

buffer(timespan, timeshift, unit) создаёт новый List элементов каждые timeshift единиц времени, и заполняет этот пакет каждым элементом, излученным исходным Observable с этого момента до тех пор, пока timespan времени не пройдёт с момента создания пакета, прежде чем излучить этот List как своё собственное излучение. Если timespan больше, чем timeshift, выпущенные пакеты будут представлять собой перекрывающиеся временные промежутки, и поэтому они могут содержать дублирующиеся элементы. Также существует версия этого варианта оператора, которая принимает Scheduler в качестве параметра и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик computation.

RxNET Buffer

В Rx.NET существует несколько вариантов Buffer. Для каждого варианта вы можете либо передать исходный Observable в качестве первого параметра, либо вызвать его как метод экземпляра исходного Observable (в этом случае вы можете опустить этот параметр):

Buffer(count)

Buffer(count)

Buffer(count) излучает непересекающиеся буферы в виде ILists, каждый из которых содержит не более count элементов из исходного Observable (последний выпущенный IList может содержать меньше, чем count элементов).

Buffer(count, skip)

Buffer(count,skip)

Buffer(count, skip) создаёт новый буфер, начиная с первого излученного элемента из исходного Observable, и каждые skip элементы после этого, и заполняет каждый буфер count элементами: начальным элементом и count-1 последующими. Эти буферы излучаются как ILists. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент), или у них могут быть пробелы (где элементы, излученные исходным Observable, не представлены ни в одном буфере).

Buffer(bufferClosingSelector)

Buffer(bufferClosingSelector)

При подписке на исходный Observable, Buffer(bufferClosingSelector) начинает собирать его излучения в IList, а также вызывает bufferClosingSelector для генерации второго Observable. Когда этот второй Observable излучает объект TBufferClosing, Buffer излучает текущий IList и повторяет этот процесс: начинает новый IList и вызывает bufferClosingSelector для создания нового Observable для мониторинга. Это будет продолжаться до завершения исходного Observable.

Buffer(bufferOpenings,bufferClosingSelector)

Buffer(bufferOpenings,bufferClosingSelector)

Buffer(bufferOpenings, bufferClosingSelector) отслеживает Observable, BufferOpenings, который излучает TBufferOpening объекты. Каждый раз, когда он наблюдает такой выпущенный элемент, он создаёт новый IList для начала сбора элементов, излучаемых исходным Observable, и передаёт TBufferOpening объект в функцию bufferClosingSelector. Эта функция возвращает Observable. Buffer отслеживает этот Observable, и когда он обнаруживает выпущенный элемент из него, он закрывает IList и излучает его как своё собственное излучение.

Buffer(timeSpan)

Buffer(timeSpan)

Buffer(timeSpan) излучает новый IList элементов периодически, каждые timeSpan единиц времени, содержащий все элементы, излученные исходным Observable с момента предыдущего излучения пакета или, в случае первого списка, с момента подписки на исходный Observable. Также существует версия этого варианта оператора, которая принимает IScheduler в качестве параметра и использует его для управления временным интервалом.

Buffer(timeSpan, count)

Buffer(timeSpan,count)

Buffer(timeSpan, count) излучает новый IList элементов для каждого count элементов, излученных исходным Observable, или, если timeSpan истек с момента последнего излучения списка, он излучает список из количества элементов, которые исходный Observable излучил за этот промежуток времени, даже если это меньше count. Также существует версия этого варианта оператора, которая принимает IScheduler в качестве параметра и использует его для управления временным интервалом.

Buffer(timeSpan, timeShift)

Buffer(timeSpan,timeShift)

Buffer(timeSpan, timeShift) создаёт новый IList элементов каждые timeShift единиц времени, и заполняет этот список каждым элементом, излученным исходным Observable с этого момента до тех пор, пока timeSpan времени не пройдёт с момента создания списка, прежде чем излучить этот IList как своё собственное излучение. Если timeSpan больше, чем timeShift, выпущенные списки будут представлять собой перекрывающиеся временные промежутки, и поэтому они могут содержать дублирующиеся элементы. Также существует версия этого варианта оператора, которая принимает IScheduler в качестве параметра и использует его для управления временным интервалом.

RxPHP bufferWithCount

RxPHP реализует этот оператор как bufferWithCount.

Проецирует каждый элемент последовательности observable в нулевой или более буферов, которые производятся на основе информации о количестве элементов.

Пример кода

//from https://github.com/ReactiveX/RxPHP/blob/master/demo/bufferWithCount/bufferWithCount.php

$source = Rx\Observable::range(1, 6)
    ->bufferWithCount(2)
    ->subscribe($stdoutObserver);
Next value: [1,2]
Next value: [3,4]
Next value: [5,6]
Complete!
//from https://github.com/ReactiveX/RxPHP/blob/master/demo/bufferWithCount/bufferWithCountAndSkip.php

$source = Rx\Observable::range(1, 6)
    ->bufferWithCount(2, 1)
    ->subscribe($stdoutObserver);
Next value: [1,2]
Next value: [2,3]
Next value: [3,4]
Next value: [4,5]
Next value: [5,6]
Next value: [6]
Complete!

RxPY buffer buffer_with_count buffer_with_time buffer_with_time_or_count pairwise

RxPY имеет несколько вариантов оператора Buffer: buffer, buffer_with_count, buffer_with_time, и buffer_with_time_or_count. Для каждого из этих вариантов существуют необязательные параметры, которые изменяют поведение оператора. Как всегда в RxPY, когда оператор может принимать более одного необязательного параметра, обязательно указывайте имя параметра в списке параметров при вызове оператора, чтобы избежать неоднозначности.

buffer(buffer_openings)

buffer(buffer_openings)

buffer(buffer_openings=boundaryObservable) отслеживает Observable, buffer_openings. Каждый раз, когда этот Observable излучает элемент, он создаёт новый массив для сбора элементов, излучаемых исходным Observable, и излучает предыдущий массив.

buffer(closing_selector)

buffer(closing_selector)

buffer(closing_selector=closingSelector) начинает собирать элементы, излучаемые исходным Observable, сразу после подписки, а также вызывает функцию closing_selector, чтобы сгенерировать второй Observable. Он отслеживает этот новый Observable и, когда он завершается или излучает элемент, излучает текущий массив, начинает новый массив для сбора элементов из исходного Observable и снова вызывает closing_selector, чтобы сгенерировать новый Observable для отслеживания, чтобы определить, когда излучить новый массив. Он повторяет этот процесс до тех пор, пока исходный Observable не завершит работу, после чего излучает последний массив.

buffer(closing_selector,buffer_closing_selector)

buffer(closing_selector=openingSelector, buffer_closing_selector=closingSelector) начинает с вызова closing_selector, чтобы получить Observable. Он отслеживает этот Observable и, всякий раз, когда он излучает элемент, buffer создаёт новый массив, начинает собирать в него элементы, которые впоследствии излучает исходный Observable, и вызывает buffer_closing_selector, чтобы получить новый Observable для управления закрытием этого массива. Когда этот новый Observable излучает элемент или завершается, buffer закрывает и излучает массив, которым управляет Observable.

buffer_with_count(count)

buffer_with_count(count)

buffer_with_count(count) излучает неперекрывающиеся буферы в виде массивов, каждый из которых содержит не более count элементов из исходного Observable (последний излучаемый массив может содержать меньше, чем count элементов).

buffer_with_count(count, skip)

buffer_with_count(count,skip)

buffer_with_count(count, skip=skip) создаёт новый буфер, начиная с первого излученного элемента от исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальный элемент и count-1 последующие. Он излучает эти буферы в виде массивов. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент), или они могут иметь пробелы (где элементы, излучаемые исходным Observable, не представлены ни в одном буфере).

buffer_with_time(timespan)

buffer_with_time(timespan)

buffer_with_time(timespan) излучает новый массив элементов периодически, каждые timespan миллисекунд, содержащий все элементы, излучаемые исходным Observable с момента предыдущего излучения буфера или, в случае первого буфера, с момента подписки на исходный Observable. Также существует версия этого варианта оператора, которая принимает параметр scheduler и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик timeout.

buffer_with_time(timespan, timeshift)

buffer_with_time(timespan,timeshift)

buffer(timespan, timeshift=timeshift) создаёт новый массив элементов каждые timeshift миллисекунд и заполняет этот массив каждым элементом, излученным исходным Observable с этого момента, пока не пройдёт timespan миллисекунд с момента создания массива, прежде чем излучить этот массив как собственное излучение. Если timespan длиннее timeshift, излучаемые массивы будут представлять перекрывающиеся временные интервалы, поэтому они могут содержать дублируемые элементы. Также существует версия этого варианта оператора, которая принимает параметр scheduler и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик timeout.

buffer_with_time_or_count(timespan, count)

buffer_with_time_or_count(timespan,count)

buffer_with_time_or_count(timespan, count) излучает новый массив элементов для каждых count элементов, излучаемых исходным Observable, или, если прошло timespan миллисекунд с момента его последнего излучения буфера, он излучает массив из того количества элементов, которые исходный Observable излучил за этот промежуток времени, даже если это меньше count. Также существует версия этого варианта оператора, которая принимает параметр scheduler и использует его для управления временным интервалом; по умолчанию этот вариант использует планировщик timeout.

Rxrb buffer_with_count buffer_with_time

Rx.rb имеет три варианта оператора Buffer:

buffer_with_count(count)

buffer_with_count(count)

buffer_with_count(count) излучает неперекрывающиеся буферы в виде массивов, каждый из которых содержит не более count элементов из исходного Observable (последний излучаемый массив может содержать меньше, чем count элементов).

buffer_with_count(count,skip)

buffer_with_count(count,skip)

buffer_with_count(count, skip=skip) создаёт новый буфер, начиная с первого излученного элемента от исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальный элемент и count-1 последующие. Он излучает эти буферы в виде массивов. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент), или они могут иметь пробелы (где элементы, излучаемые исходным Observable, не представлены ни в одном буфере).

buffer_with_time(timespan)

buffer_with_time(timespan)

buffer_with_time(timespan) излучает новый массив элементов периодически, каждые timespan миллисекунд, содержащий все элементы, излучаемые исходным Observable с момента предыдущего излучения буфера или, в случае первого буфера, с момента подписки на исходный Observable.

RxScala slidingBuffer tumblingBuffer

RxScala имеет два варианта Buffer — slidingBuffer и tumblingBuffer — каждый из которых имеет варианты с различными способами сборки буферов, которые они излучают:

slidingBuffer(count, skip)

slidingBuffer(count,skip)

slidingBuffer(count, skip) создает новый буфер, начиная с первого излученного элемента из исходного Observable, и каждые skip элементов после этого, и заполняет каждый буфер count элементами: начальный элемент и count-1 последующие. Он излучает эти буферы как Seq. В зависимости от значений count и skip эти буферы могут перекрываться (несколько буферов могут содержать один и тот же элемент) или иметь пробелы (где элементы, излучаемые исходным Observable, не представлены ни в одном буфере).

slidingBuffer(timespan, timeshift)

slidingBuffer(timespan,timeshift)

slidingBuffer(timespan, timeshift) создает новый Seq элементов каждые timeshift (a Duration period), и заполняет этот буфер каждым элементом, излучаемым исходным Observable с этого момента до тех пор, пока timespan (также Duration period) не пройдет с момента создания буфера, прежде чем излучить этот Seq в качестве собственного излучения. Если timespan длиннее, чем timeshift, излучаемые массивы будут представлять перекрывающиеся временные периоды, и поэтому они могут содержать дублирующиеся элементы. Существует также версия этого варианта оператора, который принимает Scheduler в качестве параметра и использует его для управления временным интервалом.

slidingBuffer(openings, closings)

slidingBuffer(openings,closings)

slidingBuffer(openings,closings) отслеживает openings Observable, и всякий раз, когда он излучает Opening элемент, slidingBuffer создает новый Seq, начинает собирать элементы, излучаемые исходным Observable в этот буфер, и вызывает closings для получения нового Observable для управления закрытием этого буфера. Когда этот новый Observable излучает элемент или завершается, slidingBuffer закрывается и излучает Seq, которое управляет Observable.

tumblingBuffer(count)

tumblingBuffer(count)

tumblingBuffer(count) излучает неперекрывающиеся буферы в форме Seq , каждый из которых содержит не более count элементов из исходного Observable (последний излучаемый буфер может содержать меньше, чем count элементов).

tumblingBuffer(boundary)

tumblingBuffer(boundary)

tumblingBuffer(boundary) отслеживает Observable, boundary. Каждый раз, когда этот Observable излучает элемент, он создаёт новый Seq для начала сбора элементов, излучаемых исходным Observable, и излучает предыдущий Seq. Этот вариант оператора имеет необязательный второй параметр, initialCapacity, с помощью которого вы можете указать ожидаемый размер этих буферов для повышения эффективности выделения памяти.

tumblingBuffer(timespan)

tumblingBuffer(timespan)

tumblingBuffer(timespan) излучает новый Seq элементов периодически, каждые timespan (a Duration period), содержащий все элементы, излучаемые исходным Observable с момента предыдущего излучения пакета или, в случае первого пакета, с момента подписки на исходный Observable. Этот вариант оператора имеет необязательный второй параметр, scheduler, с помощью которого можно установить Scheduler, который вы хотите использовать для управления расчетом временного интервала.

tumblingBuffer(timespan, count)

tumblingBuffer(timespan,count)

tumblingBuffer(timespan, count) излучает новый Seq элементов для каждого count элементов, излучаемых исходным Observable, или, если timespan (a Duration period) прошло с момента последнего излучения пакета, он излучает Seq содержащий столько элементов, сколько излучил исходный Observable за этот период, даже если это меньше, чем count. Этот вариант оператора имеет необязательный третий параметр, scheduler, с помощью которого вы можете установить Scheduler, который вы хотите использовать для управления расчетом временного интервала.

© ReactiveX contributors
Licensed under the Apache License 2.0.
http://reactivex.io/documentation/operators/buffer.html

Spec-Zone.ru

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