Источник Поток
Функции для создания и комбинирования потоков.
Потоки — это комбинируемые, ленивые перечисляемые (для ознакомления с перечисляемыми, см. модуль Enum). Любой перечисляемый, который генерирует элементы один за другим во время перечисления, называется потоком. Например, в Elixir 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 создаёт рецепт вычислений, которые выполняются в более поздний момент. Посмотрим ещё один пример:
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 и другие.
Обратите внимание, что функции в этом модуле гарантированно возвращают перечисляемые объекты. Поскольку перечисляемые объекты могут иметь разную структуру (структуры, анонимные функции и так далее), функции в этом модуле могут возвращать любой из этих типов, и это может измениться в любое время. Например, функция, которая сегодня возвращает анонимную функцию, может возвращать структуру в будущих версиях.
Резюме
Типы
- index()
Индекс, начинающийся с нуля.
Функции
- 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)
Сцепляет соответствующие элементы из конечного набора перечисляемых объектов в один поток кортежей.
- zip(enumerable1, enumerable2)
Ленивое объединение двух перечислимых объектов.
- zip_with(enumerables, zip_fun)
Ленивое объединение соответствующих элементов из конечного набора перечислимых объектов в новый перечислимый объект, преобразуя их с помощью функции
zip_funпо ходу работы.- zip_with(enumerable1, enumerable2, zip_fun)
Ленивое объединение соответствующих элементов из двух перечислимых объектов в новый, преобразуя их с помощью функции
zip_funпо ходу работы.
Типы
Функции
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 к enumerable и уплощает результат.
Эта функция возвращает новый поток, составленный путём добавления результата вызова 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]
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: any() Преобразует существующий поток.
Ожидает аккумулятор и функцию, которая получает два аргумента: элемент потока и обновлённый аккумулятор. Она должна вернуть кортеж, где первый элемент — новый поток (часто список) или атом :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: any() Аналогично 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: any() Преобразует существующий поток с функциями-обработчиками начала, конца и послеобработки.
После начала преобразования вызывается start_fun для вычисления начального аккумулятора. Затем для каждого элемента в перечислимом объекте вызывается функция reducer, которая получает элемент и аккумулятор, возвращая новые элементы и новый аккумулятор, как в 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]
© 2012-2024 The Elixir Team
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.17.2/Stream.html