Spec-Zone.ru › Elixir 1.18

Источник Поток

Функции для создания и композиции потоков.

Потоки являются композируемыми, ленивыми перечисляемыми (для ознакомления с перечисляемыми, см. модуль Enum). Любой перечисляемый, который генерирует элементы по одному во время перечисления, называется потоком. Например, диапазон Elixir's Range является потоком:

iex> range = 1..5
1..5
iex> Enum.map(range, &(&1 * 2))
[2, 4, 6, 8, 10]

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

iex> range = 1..3
iex> stream = Stream.map(range, &(&1 * 2))
iex> Enum.map(stream, &(&1 + 1))
[3, 5, 7]

Обратите внимание, что мы начали с диапазона, а затем создали поток, предназначенный для умножения каждого элемента в диапазоне на 2. На этом этапе вычисления не производились. Только когда вызывается Enum.map/2, мы фактически перечисляем каждый элемент в диапазоне, умножаем его на 2 и добавляем 1. Мы говорим, что функции в Stream являются ленивыми, а функции в Enum являются жадными.

Благодаря ленивости, потоки полезны при работе с большими (или даже бесконечными) коллекциями. При цепочке многих операций с Enum создаются промежуточные списки, тогда как Stream создает рецепт вычислений, которые выполняются в более поздний момент. Затем, когда поток потребляется позже, чаще всего с помощью функции в модуле Enum, поток будет выдавать свои элементы по одному.

Давайте рассмотрим другой пример:

1..3
|> Enum.map(&IO.inspect(&1))
|> Enum.map(&(&1 * 2))
|> Enum.map(&IO.inspect(&1))
1
2
3
2
4
6
#=> [2, 4, 6]

Обратите внимание, что мы сначала вывели каждый элемент в списке, затем умножили каждый элемент на 2 и, наконец, вывели каждое новое значение. В этом примере список был перечислен три раза. Давайте посмотрим пример с потоками:

stream = 1..3
|> Stream.map(&IO.inspect(&1))
|> Stream.map(&(&1 * 2))
|> Stream.map(&IO.inspect(&1))
Enum.to_list(stream)
1
2
2
4
3
6
#=> [2, 4, 6]

Хотя конечный результат тот же, порядок, в котором элементы выводились на экран, изменился! С потоками мы выводим первый элемент, а затем выводим его удвоенное значение. В этом примере список был перечислен только один раз!

Вот что мы имели в виду, когда ранее говорили, что потоки являются композируемыми, ленивыми перечисляемыми. Обратите внимание, что мы можем вызвать Stream.map/2 несколько раз, эффективно комбинируя потоки и сохраняя их ленивыми. Вычисления выполняются только при вызове функции из модуля Enum.

Как и с Enum, функции в этом модуле работают за линейное время. Это означает, что время, затрачиваемое на выполнение операции, увеличивается с той же скоростью, что и длина списка. Это ожидается при операциях, таких как Stream.map/2. Ведь если мы хотим пройти по каждому элементу в потоке, чем длиннее поток, тем больше элементов нам нужно пройти, и тем дольше это займет.

Создание потоков

Существует много функций в стандартной библиотеке Elixir, которые возвращают потоки, некоторые примеры:

  • IO.stream/2 - потоки входных строк, по одной строке
  • URI.query_decoder/1 - декодирует строку запроса, пару за парой

Этот модуль также предоставляет множество удобных функций для создания потоков, таких как Stream.cycle/1, Stream.unfold/2, Stream.resource/3 и другие.

Не проверяйте структуры Stream

Хотя некоторые функции в этом модуле могут возвращать структуру Stream, вы никогда не должны явно проверять структуру Stream, так как потоки могут иметь различную форму, такую как IO.Stream, File.Stream, или даже диапазоны Range.

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

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

Резюме

Типы

acc()
default()
element()
index()

Индекс с нулевым основанием.

timer()

Функции

chunk_by(enum, fun)

Разбивает enum на куски, буферизуя элементы, для которых fun возвращает одинаковое значение.

chunk_every(enum, count)

Сокращение для chunk_every(enum, count, count).

chunk_every(enum, count, step, leftover \\ [])

Потоковый ввод перечисляемого объекта кусками, содержащими по count элементов каждый, где каждый новый кусок начинается через step элементов в перечисляемом объекте.

chunk_while(enum, acc, chunk_fun, after_fun)

Разбивает enum на куски с подробным управлением, когда каждый кусок отправляется.

concat(enumerables)

Создаёт поток, который перечисляет каждый перечисляемый объект в перечисляемом объекте.

concat(first, second)

Создаёт поток, который перечисляет первый аргумент, а затем второй.

cycle(enumerable)

Создаёт поток, который циклически проходит по заданному перечисляемому объекту бесконечно.

dedup(enum)

Создаёт поток, который отправляет элементы только в том случае, если они отличаются от последнего отправленного элемента.

dedup_by(enum, fun)

Создаёт поток, который отправляет элементы только в том случае, если результат вызова fun на элементе отличается от (сохранённого) результата вызова fun на последнем отправленном элементе.

drop(enum, n)

Лениво пропускает следующие n элементов из перечисляемого объекта.

drop_every(enum, nth)

Создаёт поток, который пропускает каждый nth элемент из перечисляемого объекта.

drop_while(enum, fun)

Лениво пропускает элементы перечисляемого объекта, пока заданная функция возвращает истинное значение.

duplicate(value, n)

Дублирует заданный элемент n раз в потоке.

each(enum, fun)

Выполняет заданную функцию для каждого элемента.

filter(enum, fun)

Создаёт поток, который фильтрует элементы в соответствии с заданной функцией при перечислении.

flat_map(enum, mapper)

Применяет заданную fun к enumerable и сплющивает результат.

from_index(fun_or_offset \\ 0)

Создаёт поток из индекса, начиная с смещения или заданного функцией.

intersperse(enumerable, intersperse_element)

Лениво вставляет intersperse_element между каждым элементом перечисления.

interval(n)

Создаёт поток, который отправляет значение после заданного периода n в миллисекундах.

into(enum, collectable, transform \\ fn x -> x end)

Вставляет значения потока в заданный коллектор как побочный эффект.

iterate(start_value, next_fun)

Отправляет последовательность значений, начиная с start_value.

map(enum, fun)

Создаёт поток, который применит заданную функцию к перечислению.

map_every(enum, nth, fun)

Создаёт поток, который применит заданную функцию к каждому nth элементу из перечисляемого объекта.

reject(enum, fun)

Создаёт поток, который отклонит элементы в соответствии с заданной функцией при перечислении.

repeatedly(generator_fun)

Возвращает поток, сгенерированный путём вызова generator_fun многократно.

resource(start_fun, next_fun, after_fun)

Отправляет последовательность значений для данного ресурса.

run(stream)

Выполняет заданный поток.

scan(enum, fun)

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

scan(enum, acc, fun)

Создаёт поток, который применяет заданную функцию к каждому элементу, отправляет результат и использует тот же результат в качестве аккумулятора для следующего вычисления. Использует указанное acc в качестве начального значения.

take(enum, count)

Лениво берёт следующие count элементов из перечисляемого объекта и останавливает перечисление.

take_every(enum, nth)

Создаёт поток, который берёт каждый nth элемент из перечисляемого объекта.

take_while(enum, fun)

Лениво берёт элементы перечисляемого объекта, пока заданная функция возвращает истинное значение.

timer(n)

Создаёт поток, который отправляет единственное значение после n миллисекунд.

transform(enum, acc, reducer)

Преобразует существующий поток.

transform(enum, start_fun, reducer, after_fun)

Аналогично Stream.transform/5, за исключением того, что last_fun не предоставляется.

transform(enum, start_fun, reducer, last_fun, after_fun)

Преобразует существующий поток с функциями-обработчиками начала, конца и после выполнения.

unfold(next_acc, next_fun)

Отправляет последовательность значений для данного аккумулятора.

uniq(enum)

Создаёт поток, который отправляет только уникальные элементы.

uniq_by(enum, fun)

Создаёт поток, который отправляет только уникальные элементы, удаляя элементы, для которых функция fun возвращала дублирующие элементы.

with_index(enum, fun_or_offset \\ 0)

Создаёт поток, где каждый элемент в перечисляемом объекте будет заключён в кортеж вместе с его индексом или в соответствии с заданной функцией.

zip(enumerables)

Сцепляет соответствующие элементы из конечного набора перечисляемых объектов в один поток кортежей.

END_OF_DOCUMENT_MARKER
zip(enumerable1, enumerable2)

Объединяет два перечислителя вместе, лениво.

zip_with(enumerables, zip_fun)

Лениво объединяет соответствующие элементы из конечного набора перечислителей в новый перечислитель, преобразуя их с помощью функции zip_fun по мере необходимости.

zip_with(enumerable1, enumerable2, zip_fun)

Лениво объединяет соответствующие элементы из двух перечислителей в новый, преобразуя их с помощью функции zip_fun по мере необходимости.

Типы

acc()Source

@type acc() :: any()

default()Source

@type default() :: any()

element()Source

@type element() :: any()

index()Source

@type index() :: non_neg_integer()

Индекс, отсчитываемый с нуля.

timer()Source

@type timer() :: non_neg_integer() | :infinity

Функции

chunk_by(enum, fun)Source

@spec chunk_by(Enumerable.t(), (element() -> any())) :: Enumerable.t()

Разбивает enum на куски, буферизуя элементы, для которых fun возвращает одинаковое значение.

Элементы выводятся только тогда, когда fun возвращает новое значение или enum завершается.

Примеры

iex> stream = Stream.chunk_by([1, 2, 2, 3, 4, 4, 6, 7, 7], &(rem(&1, 2) == 1))
iex> Enum.to_list(stream)
[[1], [2, 2], [3], [4, 4, 6], [7, 7]]

chunk_every(enum, count)Source

@spec chunk_every(Enumerable.t(), pos_integer()) :: Enumerable.t()

Сокращение для chunk_every(enum, count, count).

chunk_every(enum, count, step, leftover \\ [])Source

@spec chunk_every(
  Enumerable.t(),
  pos_integer(),
  pos_integer(),
  Enumerable.t() | :discard
) ::
  Enumerable.t()

Выводит перечисляемый объект частями, содержащими по count элементов каждая, где каждый новый кусок начинается через step элементов в перечисляемом объекте.

step — необязательный параметр, и если его не указать, он по умолчанию равен count, то есть куски не перекрываются. Разбиение на куски прекращается, как только коллекция закончится или когда будет выведен неполный кусок.

Если в последнем куске недостаточно count элементов для заполнения куска, элементы берутся из leftover для заполнения куска. Если в leftover недостаточно элементов для заполнения куска, возвращается частичный кусок с менее чем count элементами.

Если :discard указан в leftover, последний кусок отбрасывается, если только он не содержит ровно count элементов.

Примеры

iex> Stream.chunk_every([1, 2, 3, 4, 5, 6], 2) |> Enum.to_list()
[[1, 2], [3, 4], [5, 6]]

iex> Stream.chunk_every([1, 2, 3, 4, 5, 6], 3, 2, :discard) |> Enum.to_list()
[[1, 2, 3], [3, 4, 5]]

iex> Stream.chunk_every([1, 2, 3, 4, 5, 6], 3, 2, [7]) |> Enum.to_list()
[[1, 2, 3], [3, 4, 5], [5, 6, 7]]

iex> Stream.chunk_every([1, 2, 3, 4, 5, 6], 3, 3, []) |> Enum.to_list()
[[1, 2, 3], [4, 5, 6]]

iex> Stream.chunk_every([1, 2, 3, 4], 3, 3, Stream.cycle([0])) |> Enum.to_list()
[[1, 2, 3], [4, 0, 0]]

chunk_while(enum, acc, chunk_fun, after_fun)Source

@spec chunk_while(
  Enumerable.t(),
  acc(),
  (element(), acc() -> {:cont, chunk, acc()} | {:cont, acc()} | {:halt, acc()}),
  (acc() -> {:cont, chunk, acc()} | {:cont, acc()})
) :: Enumerable.t()
when chunk: any()

Разбивает enum с тонким управлением, когда каждый кусок выводится.

chunk_fun получает текущий элемент и накопитель и должен вернуть {:cont, element, acc}, чтобы вывести данный кусок и продолжить с накопителем, или {:cont, acc}, чтобы не выводить кусок и продолжить с возвращённым накопителем.

after_fun вызывается при завершении итерации и также должен вернуть {:cont, element, acc} или {:cont, acc}.

Примеры

iex> chunk_fun = fn element, acc ->
...>   if rem(element, 2) == 0 do
...>     {:cont, Enum.reverse([element | acc]), []}
...>   else
...>     {:cont, [element | acc]}
...>   end
...> end
iex> after_fun = fn
...>   [] -> {:cont, []}
...>   acc -> {:cont, Enum.reverse(acc), []}
...> end
iex> stream = Stream.chunk_while(1..10, [], chunk_fun, after_fun)
iex> Enum.to_list(stream)
[[1, 2], [3, 4], [5, 6], [7, 8], [9, 10]]

concat(enumerables)Source

@spec concat(Enumerable.t()) :: Enumerable.t()

Создаёт поток, который перечисляет каждый перечисляемый объект в перечисляемом объекте.

Примеры

iex> stream = Stream.concat([1..3, 4..6, 7..9])
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5, 6, 7, 8, 9]

concat(first, second)Source

@spec concat(Enumerable.t(), Enumerable.t()) :: Enumerable.t()

Создаёт поток, который перечисляет первый аргумент, а затем второй.

Примеры

iex> stream = Stream.concat(1..3, 4..6)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5, 6]

iex> stream1 = Stream.cycle([1, 2, 3])
iex> stream2 = Stream.cycle([4, 5, 6])
iex> stream = Stream.concat(stream1, stream2)
iex> Enum.take(stream, 6)
[1, 2, 3, 1, 2, 3]

cycle(enumerable)Source

@spec cycle(Enumerable.t()) :: Enumerable.t()

Создаёт поток, который циклически перечисляет заданный перечисляемый объект бесконечно.

Примеры

iex> stream = Stream.cycle([1, 2, 3])
iex> Enum.take(stream, 5)
[1, 2, 3, 1, 2]

dedup(enum)Source

@spec dedup(Enumerable.t()) :: Enumerable.t()

Создаёт поток, который выводит элементы только если они отличаются от последнего выведенного элемента.

Для этой функции необходимо хранить только последний выведенный элемент.

Элементы сравниваются с помощью ===/2.

Примеры

iex> Stream.dedup([1, 2, 3, 3, 2, 1]) |> Enum.to_list()
[1, 2, 3, 2, 1]

dedup_by(enum, fun)Source

@spec dedup_by(Enumerable.t(), (element() -> term())) :: Enumerable.t()

Создаёт поток, который выводит элементы только если результат вызова fun на элементе отличается от результата вызова fun на последнем выведенном элементе.

Примеры

iex> Stream.dedup_by([{1, :x}, {2, :y}, {2, :z}, {1, :x}], fn {x, _} -> x end) |> Enum.to_list()
[{1, :x}, {2, :y}, {1, :x}]

drop(enum, n)Source

@spec drop(Enumerable.t(), integer()) :: Enumerable.t()

Лениво пропускает следующие n элементов из перечисляемого объекта.

Если задан отрицательный n, он пропустит последние n элементов из коллекции. Обратите внимание, что механизм, с помощью которого это реализовано, откладывает вывод любого элемента до тех пор, пока не будут выведены n дополнительных элементов из перечисления.

Примеры

iex> stream = Stream.drop(1..10, 5)
iex> Enum.to_list(stream)
[6, 7, 8, 9, 10]

iex> stream = Stream.drop(1..10, -5)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

drop_every(enum, nth)Source

@spec drop_every(Enumerable.t(), non_neg_integer()) :: Enumerable.t()

Создаёт поток, который пропускает каждый nth элемент из перечисляемого объекта.

Первый элемент всегда пропускается, если nth не равно 0.

nth должно быть неотрицательным целым числом.

Примеры

iex> stream = Stream.drop_every(1..10, 2)
iex> Enum.to_list(stream)
[2, 4, 6, 8, 10]

iex> stream = Stream.drop_every(1..1000, 1)
iex> Enum.to_list(stream)
[]

iex> stream = Stream.drop_every([1, 2, 3, 4, 5], 0)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

drop_while(enum, fun)Source

@spec drop_while(Enumerable.t(), (element() -> as_boolean(term()))) :: Enumerable.t()

Лениво пропускает элементы перечисляемого объекта, пока заданная функция возвращает истинное значение.

Примеры

iex> stream = Stream.drop_while(1..10, &(&1 <= 5))
iex> Enum.to_list(stream)
[6, 7, 8, 9, 10]

duplicate(value, n)Source

@spec duplicate(any(), non_neg_integer()) :: Enumerable.t()

Дублирует заданный элемент n раз в потоке.

n — целое число, большее или равное 0.

Если n равно 0, возвращается пустой поток.

Примеры

iex> stream = Stream.duplicate("hello", 0)
iex> Enum.to_list(stream)
[]

iex> stream = Stream.duplicate("hi", 1)
iex> Enum.to_list(stream)
["hi"]

iex> stream = Stream.duplicate("bye", 2)
iex> Enum.to_list(stream)
["bye", "bye"]

iex> stream = Stream.duplicate([1, 2], 3)
iex> Enum.to_list(stream)
[[1, 2], [1, 2], [1, 2]]

each(enum, fun)Source

@spec each(Enumerable.t(), (element() -> term())) :: Enumerable.t()

Выполняет заданную функцию для каждого элемента.

Значения в потоке не изменяются, поэтому эта функция полезна для добавления побочных эффектов (например, вывода на печать) в поток. См. map/2, если требуется создание другого потока.

Примеры

iex> stream = Stream.each([1, 2, 3], fn x -> send(self(), x) end)
iex> Enum.to_list(stream)
iex> receive do: (x when is_integer(x) -> x)
1
iex> receive do: (x when is_integer(x) -> x)
2
iex> receive do: (x when is_integer(x) -> x)
3

filter(enum, fun)Source

@spec filter(Enumerable.t(), (element() -> as_boolean(term()))) :: Enumerable.t()

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

Примеры

iex> stream = Stream.filter([1, 2, 3], fn x -> rem(x, 2) == 0 end)
iex> Enum.to_list(stream)
[2]

flat_map(enum, mapper)Source

@spec flat_map(Enumerable.t(), (element() -> Enumerable.t())) :: Enumerable.t()

Применяет заданную функцию-преобразователь к данному fun и сглаживает результат.

Эта функция возвращает новый поток, созданный путём добавления результата вызова fun для каждого элемента из enumerable.

Примеры

iex> stream = Stream.flat_map([1, 2, 3], fn x -> [x, x * 2] end)
iex> Enum.to_list(stream)
[1, 2, 2, 4, 3, 6]

iex> stream = Stream.flat_map([1, 2, 3], fn x -> [[x]] end)
iex> Enum.to_list(stream)
[[1], [2], [3]]

from_index(fun_or_offset \\ 0)Source

@spec from_index(integer()) :: Enumerable.t(integer())
@spec from_index((integer() -> return_value)) :: Enumerable.t(return_value)
when return_value: term()

Создаёт поток из индекса, либо начиная с смещения, либо используя функцию.

Может принять функцию или целое смещение.

Если задано offset, он выведет элементы со смещения.

Если задана function, он вызовет функцию с элементами со смещения.

Примеры

iex> Stream.from_index() |> Enum.take(3)
[0, 1, 2]

iex> Stream.from_index(1) |> Enum.take(3)
[1, 2, 3]

iex> Stream.from_index(fn x -> x * 10 end) |> Enum.take(3)
[0, 10, 20]
END_OF_DOCUMENT_MARKER

intersperse(enumerable, intersperse_element)Source

@spec intersperse(Enumerable.t(), any()) :: Enumerable.t()

Лениво вставляет intersperse_element между каждым элементом перечисления.

Примеры

iex> Stream.intersperse([1, 2, 3], 0) |> Enum.to_list()
[1, 0, 2, 0, 3]

iex> Stream.intersperse([1], 0) |> Enum.to_list()
[1]

iex> Stream.intersperse([], 0) |> Enum.to_list()
[]

interval(n)Source

@spec interval(timer()) :: Enumerable.t()

Создаёт поток, который излучает значение через заданный период n в миллисекундах.

Излучаемые значения — это возрастающий счётчик, начинающийся с 0. Эта операция блокирует вызывающую функцию заданным интервалом каждый раз, когда излучается новый элемент.

Не используйте эту функцию для генерации последовательности чисел. Если блокировка вызывающего процесса не нужна, используйте Stream.iterate(0, & &1 + 1) вместо неё.

Примеры

iex> Stream.interval(10) |> Enum.take(10)
[0, 1, 2, 3, 4, 5, 6, 7, 8, 9]

into(enum, collectable, transform \\ fn x -> x end)Source

@spec into(Enumerable.t(), Collectable.t(), (term() -> term())) :: Enumerable.t()

Вводит значения потока в задаваемый коллектор как побочный эффект.

Эта функция часто используется с run/1, так как любая оценка откладывается до выполнения потока. См. run/1 для примера.

iterate(start_value, next_fun)Source

@spec iterate(element(), (element() -> element())) :: Enumerable.t()

Излучает последовательность значений, начиная с start_value.

Последовательные значения генерируются путём вызова next_fun над предыдущим значением.

Примеры

iex> Stream.iterate(1, &(&1 * 2)) |> Enum.take(5)
[1, 2, 4, 8, 16]

map(enum, fun)Source

@spec map(Enumerable.t(), (element() -> any())) :: Enumerable.t()

Создаёт поток, который будет применять заданную функцию к перечислению.

Примеры

iex> stream = Stream.map([1, 2, 3], fn x -> x * 2 end)
iex> Enum.to_list(stream)
[2, 4, 6]

map_every(enum, nth, fun)Source

@spec map_every(Enumerable.t(), non_neg_integer(), (element() -> any())) ::
  Enumerable.t()

Создаёт поток, который будет применять заданную функцию к каждому nth элементу из перечисления.

Первый элемент всегда передаётся в заданную функцию.

nth должно быть целым неотрицательным числом.

Примеры

iex> stream = Stream.map_every(1..10, 2, fn x -> x * 2 end)
iex> Enum.to_list(stream)
[2, 2, 6, 4, 10, 6, 14, 8, 18, 10]

iex> stream = Stream.map_every([1, 2, 3, 4, 5], 1, fn x -> x * 2 end)
iex> Enum.to_list(stream)
[2, 4, 6, 8, 10]

iex> stream = Stream.map_every(1..5, 0, fn x -> x * 2 end)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

reject(enum, fun)Source

@spec reject(Enumerable.t(), (element() -> as_boolean(term()))) :: Enumerable.t()

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

Примеры

iex> stream = Stream.reject([1, 2, 3], fn x -> rem(x, 2) == 0 end)
iex> Enum.to_list(stream)
[1, 3]

repeatedly(generator_fun)Source

@spec repeatedly((-> element())) :: Enumerable.t()

Возвращает поток, генерируемый путём многократного вызова generator_fun.

Примеры

# Although not necessary, let's seed the random algorithm
iex> :rand.seed(:exsss, {1, 2, 3})
iex> Stream.repeatedly(&:rand.uniform/0) |> Enum.take(3)
[0.5455598952593053, 0.6039309974353404, 0.6684893034823949]

resource(start_fun, next_fun, after_fun)Source

@spec resource((-> acc()), (acc() -> {[element()], acc()} | {:halt, acc()}), (acc() ->
                                                                          term())) ::
  Enumerable.t()

Излучает последовательность значений для заданного ресурса.

Аналогично transform/3, но начальное накопленное значение вычисляется лениво через start_fun и выполняется after_fun в конце перечисления (в случае успеха и неудачи).

Последовательные значения генерируются путём вызова next_fun с предыдущим аккумулятором (начальное значение — результат, возвращаемый start_fun). Она должна возвращать кортеж, содержащий список излучаемых элементов и следующий аккумулятор. Перечисление завершается, если возвращает {:halt, acc}.

Как следует из названия функции, эта функция полезна для потоковой передачи значений из ресурсов.

Примеры

Stream.resource(
  fn -> File.open!("sample") end,
  fn file ->
    case IO.read(file, :line) do
      data when is_binary(data) -> {[data], file}
      _ -> {:halt, file}
    end
  end,
  fn file -> File.close(file) end
)

iex> Stream.resource(
...>  fn ->
...>    {:ok, pid} = StringIO.open("string")
...>    pid
...>  end,
...>  fn pid ->
...>    case IO.getn(pid, "", 1) do
...>      :eof -> {:halt, pid}
...>      char -> {[char], pid}
...>    end
...>  end,
...>  fn pid -> StringIO.close(pid) end
...> ) |> Enum.to_list()
["s", "t", "r", "i", "n", "g"]

run(stream)Source

@spec run(Enumerable.t()) :: :ok

Выполняет заданный поток.

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

Примеры

Открыть файл, заменить все # на % и передать в другой файл, не загружая весь файл в память:

File.stream!("/path/to/file")
|> Stream.map(&String.replace(&1, "#", "%"))
|> Stream.into(File.stream!("/path/to/other/file"))
|> Stream.run()

Никаких вычислений не будет выполнено, пока мы не вызовем одну из функций Enum или run/1.

scan(enum, fun)Source

@spec scan(Enumerable.t(), (element(), acc() -> any())) :: Enumerable.t()

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

Примеры

iex> stream = Stream.scan(1..5, &(&1 + &2))
iex> Enum.to_list(stream)
[1, 3, 6, 10, 15]

scan(enum, acc, fun)Source

@spec scan(Enumerable.t(), acc(), (element(), acc() -> any())) :: Enumerable.t()

Создаёт поток, который применяет заданную функцию к каждому элементу, излучает результат и использует тот же результат в качестве аккумулятора для следующего вычисления. Использует заданное acc в качестве начального значения.

Примеры

iex> stream = Stream.scan(1..5, 0, &(&1 + &2))
iex> Enum.to_list(stream)
[1, 3, 6, 10, 15]

take(enum, count)Source

@spec take(Enumerable.t(), integer()) :: Enumerable.t()

Лениво берёт следующие count элементы из перечисления и останавливает перечисление.

Если задано отрицательное count, будут взяты последние count значения. В таком случае коллекция полностью перечисляется, сохраняя до 2 * count элементов в памяти. Как только конец коллекции достигнут, будут выполнены последние count элементы. Поэтому использование отрицательного count на бесконечном наборе никогда не вернёт результат.

Примеры

iex> stream = Stream.take(1..100, 5)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

iex> stream = Stream.take(1..100, -5)
iex> Enum.to_list(stream)
[96, 97, 98, 99, 100]

iex> stream = Stream.cycle([1, 2, 3]) |> Stream.take(5)
iex> Enum.to_list(stream)
[1, 2, 3, 1, 2]

take_every(enum, nth)Source

@spec take_every(Enumerable.t(), non_neg_integer()) :: Enumerable.t()

Создаёт поток, который берёт каждый nth элемент из перечисления.

Первый элемент всегда включается, если nth не равен 0.

nth должно быть целым неотрицательным числом.

Примеры

iex> stream = Stream.take_every(1..10, 2)
iex> Enum.to_list(stream)
[1, 3, 5, 7, 9]

iex> stream = Stream.take_every([1, 2, 3, 4, 5], 1)
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

iex> stream = Stream.take_every(1..1000, 0)
iex> Enum.to_list(stream)
[]

take_while(enum, fun)Source

@spec take_while(Enumerable.t(), (element() -> as_boolean(term()))) :: Enumerable.t()

Лениво берёт элементы из перечисления, пока заданная функция возвращает истинное значение.

Примеры

iex> stream = Stream.take_while(1..100, &(&1 <= 5))
iex> Enum.to_list(stream)
[1, 2, 3, 4, 5]

timer(n)Source

@spec timer(timer()) :: Enumerable.t()

Создаёт поток, который излучает единственное значение через n миллисекунд.

Излучаемое значение — 0. Эта операция блокирует вызывающую функцию заданное время, пока элемент не будет излучен.

Примеры

iex> Stream.timer(10) |> Enum.to_list()
[0]

transform(enum, acc, reducer)Source

@spec transform(Enumerable.t(), acc, fun) :: Enumerable.t()
when fun: (element(), acc -> {Enumerable.t(), acc} | {:halt, acc}), acc: term()

Преобразует существующий поток.

Ожидает аккумулятора и функции, которая получает два аргумента: элемент потока и обновлённый аккмулятор. Она должна возвращать кортеж, где первый элемент — новый поток (часто список) или атом :halt, а второй элемент — аккмулятор, который будет использоваться для следующего элемента.

Примечание: эта функция эквивалентна Enum.flat_map_reduce/3, за исключением того, что эта функция не возвращает аккмулятор после обработки потока.

Примеры

Stream.transform/3 полезна, так как может использоваться в качестве основы для реализации многих функций, определённых в этом модуле. Например, мы можем реализовать Stream.take(enum, n) следующим образом:

iex> enum = 1001..9999
iex> n = 3
iex> stream = Stream.transform(enum, 0, fn i, acc ->
...>   if acc < n, do: {[i], acc + 1}, else: {:halt, acc}
...> end)
iex> Enum.to_list(stream)
[1001, 1002, 1003]

Stream.transform/5 ещё более обобщает эту функцию, позволяя оборачивать ресурсы.

transform(enum, start_fun, reducer, after_fun)Source

@spec transform(Enumerable.t(), start_fun, reducer, after_fun) :: Enumerable.t()
when start_fun: (-> acc),
     reducer: (element(), acc -> {Enumerable.t(), acc} | {:halt, acc}),
     after_fun: (acc -> term()),
     acc: term()

Аналогично Stream.transform/5, за исключением того, что last_fun не предоставляется.

Эту функцию можно рассматривать как комбинацию Stream.resource/3 и Stream.transform/3.

transform(enum, start_fun, reducer, last_fun, after_fun)Source

@spec transform(Enumerable.t(), start_fun, reducer, last_fun, after_fun) ::
  Enumerable.t()
when start_fun: (-> acc),
     reducer: (element(), acc -> {Enumerable.t(), acc} | {:halt, acc}),
     last_fun: (acc -> {Enumerable.t(), acc} | {:halt, acc}),
     after_fun: (acc -> term()),
     acc: term()

Преобразует существующий поток с функциями-обработчиками start, last и after.

После начала преобразования вызывается start_fun для вычисления начального аккумулятора. Затем для каждого элемента в перечислимом объекте вызывается функция reducer с элементом и аккумулятором, возвращая новые элементы и новый аккумулятор, как в Stream.transform/3.

После обработки коллекции вызывается last_fun с аккумулятором для выдачи оставшихся элементов. Затем вызывается after_fun для закрытия любого ресурса, но не выдачи новых элементов. last_fun вызывается только в том случае, если перечислимый объект завершается успешно (либо потому, что он завершен, либо он сам остановился). after_fun всегда вызывается, поэтому after_fun должен использоваться для закрытия ресурсов.

unfold(next_acc, next_fun)Source

@spec unfold(acc(), (acc() -> {element(), acc()} | nil)) :: Enumerable.t()

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

Последовательные значения генерируются путем вызова next_fun с предыдущим аккумулятором, и он должен возвращать кортеж с текущим значением и следующим аккумулятором. Перечисление завершается, если возвращается nil.

Примеры

Для создания потока, который считает вниз и останавливается перед нулём:

iex> Stream.unfold(5, fn
...>   0 -> nil
...>   n -> {n, n - 1}
...> end) |> Enum.to_list()
[5, 4, 3, 2, 1]

Если next_fun никогда не возвращает nil, возвращаемый поток является бесконечным:

iex> Stream.unfold(0, fn
...>   n -> {n, n + 1}
...> end) |> Enum.take(10)
[0, 1, 2, 3, 4, 5, 6, 7, 8, 9]

iex> Stream.unfold(1, fn
...>   n -> {n, n * 2}
...> end) |> Enum.take(10)
[1, 2, 4, 8, 16, 32, 64, 128, 256, 512]

uniq(enum)Source

@spec uniq(Enumerable.t()) :: Enumerable.t()

Создаёт поток, который выдает элементы только в случае их уникальности.

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

Примеры

iex> Stream.uniq([1, 2, 3, 3, 2, 1]) |> Enum.to_list()
[1, 2, 3]

uniq_by(enum, fun)Source

@spec uniq_by(Enumerable.t(), (element() -> term())) :: Enumerable.t()

Создаёт поток, который выдает элементы только в случае их уникальности, удаляя элементы, для которых функция fun возвращала дублирующие элементы.

Функция fun сопоставляет каждый элемент с термином, который используется для определения того, являются ли два элемента дубликатами.

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

Пример

iex> Stream.uniq_by([{1, :x}, {2, :y}, {1, :z}], fn {x, _} -> x end) |> Enum.to_list()
[{1, :x}, {2, :y}]

iex> Stream.uniq_by([a: {:tea, 2}, b: {:tea, 2}, c: {:coffee, 1}], fn {_, y} -> y end) |> Enum.to_list()
[a: {:tea, 2}, c: {:coffee, 1}]

with_index(enum, fun_or_offset \\ 0)Source

@spec with_index(Enumerable.t(), integer()) :: Enumerable.t({element(), integer()})
@spec with_index(Enumerable.t(), (element(), index() -> return_value)) ::
  Enumerable.t(return_value)
when return_value: term()

Создаёт поток, где каждый элемент в перечислимом объекте будет обернут в кортеж вместе с его индексом или в соответствии с заданной функцией.

Может принимать функцию или целое число смещения.

Если задано offset, индексирование будет происходить со значения смещения, а не с нуля.

Если задана function, индексирование будет происходить путём вызова функции для каждого элемента и индекса (отсчёт с нуля) перечислимого объекта.

Примеры

iex> stream = Stream.with_index([1, 2, 3])
iex> Enum.to_list(stream)
[{1, 0}, {2, 1}, {3, 2}]

iex> stream = Stream.with_index([1, 2, 3], 3)
iex> Enum.to_list(stream)
[{1, 3}, {2, 4}, {3, 5}]

iex> stream = Stream.with_index([1, 2, 3], fn x, index -> x + index end)
iex> Enum.to_list(stream)
[1, 3, 5]

zip(enumerables)Source

@spec zip(enumerables) :: Enumerable.t()
when enumerables: [Enumerable.t()] | Enumerable.t()

Объединяет соответствующие элементы из конечного набора перечислимых объектов в один поток кортежей.

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

Примеры

iex> concat = Stream.concat(1..3, 4..6)
iex> cycle = Stream.cycle(["foo", "bar", "baz"])
iex> Stream.zip([concat, [:a, :b, :c], cycle]) |> Enum.to_list()
[{1, :a, "foo"}, {2, :b, "bar"}, {3, :c, "baz"}]

zip(enumerable1, enumerable2)Source

@spec zip(Enumerable.t(), Enumerable.t()) :: Enumerable.t()

Объединяет два перечислимых объекта вместе лениво.

Так как список кортежей с двумя элементами и атомами в качестве первого элемента кортежа является списком ключевых слов (Keyword), объединение первого Stream атомов со вторым Stream любого типа создаёт Stream, который генерирует список ключевых слов.

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

Примеры

iex> concat = Stream.concat(1..3, 4..6)
iex> cycle = Stream.cycle([:a, :b, :c])
iex> Stream.zip(concat, cycle) |> Enum.to_list()
[{1, :a}, {2, :b}, {3, :c}, {4, :a}, {5, :b}, {6, :c}]
iex> Stream.zip(cycle, concat) |> Enum.to_list()
[a: 1, b: 2, c: 3, a: 4, b: 5, c: 6]

zip_with(enumerables, zip_fun)Source

@spec zip_with(enumerables, (Enumerable.t() -> term())) :: Enumerable.t()
when enumerables: [Enumerable.t()] | Enumerable.t()

Лениво объединяет соответствующие элементы из конечного набора перечислимых объектов в новый перечислимый объект, преобразуя их с помощью функции zip_fun по мере необходимости.

Первый элемент из каждого из перечислимых объектов в enumerables помещается в список, который затем передаётся функции zip_fun с одним аргументом. Затем, вторые элементы из каждого из перечислимых объектов помещаются в список и передаются zip_fun, и так далее, пока не завершит работу любой из перечислимых объектов в enumerables.

Возвращает новый перечислимый объект с результатами вызова zip_fun.

Примеры

iex> concat = Stream.concat(1..3, 4..6)
iex> Stream.zip_with([concat, concat], fn [a, b] -> a + b end) |> Enum.to_list()
[2, 4, 6, 8, 10, 12]

iex> concat = Stream.concat(1..3, 4..6)
iex> Stream.zip_with([concat, concat, 1..3], fn [a, b, c] -> a + b + c end) |> Enum.to_list()
[3, 6, 9]

zip_with(enumerable1, enumerable2, zip_fun)Source

@spec zip_with(Enumerable.t(), Enumerable.t(), (term(), term() -> term())) ::
  Enumerable.t()

Лениво объединяет соответствующие элементы из двух перечислимых объектов в новый, преобразуя их с помощью функции zip_fun по мере необходимости.

Функция zip_fun будет вызвана с первым элементом из enumerable1 и первым элементом из enumerable2, затем со вторым элементом каждого и так далее, пока не завершит работу один из перечислимых объектов.

Примеры

iex> concat = Stream.concat(1..3, 4..6)
iex> Stream.zip_with(concat, concat, fn a, b -> a + b end) |> Enum.to_list()
[2, 4, 6, 8, 10, 12]

Скачать версию ePub

Создано с помощью ExDoc (v0.36.1) для языка программирования Elixir

© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.18.1/Stream.html

Spec-Zone.ru

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