Spec-Zone.ru › Elixir 1.6

Потоки

Модуль для создания и комбинирования потоков.

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

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.

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

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

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

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

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

Резюме

Типы

acc()
default()
element()
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)

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

each(enum, fun)

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

filter(enum, fun)

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

flat_map(enum, mapper)

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

intersperse(enumerable, intersperse_element)

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

interval(n)

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

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

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

iterate(start_value, next_fun)

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

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)

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

unfold(next_acc, next_fun)

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

uniq(enum)

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

uniq_by(enum, fun)

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

with_index(enum, offset \\ 0)

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

zip(enumerables)

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

zip(left, right)

Лениво сжимает два набора вместе

Типы

acc()

acc() :: any()

default()

default() :: any()

element()

element() :: any()

index()

index() :: non_neg_integer()

Функции

chunk_by(enum, fun)

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)

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

Является сокращением для chunk_every(enum, count, count).

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

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]]

chunk_while(enum, acc, chunk_fun, after_fun)

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 i, acc ->
...>   if rem(i, 2) == 0 do
...>     {:cont, Enum.reverse([i | acc]), []}
...>   else
...>     {:cont, [i | 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)

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)

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)

cycle(Enumerable.t()) :: Enumerable.t()

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

Примеры

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

dedup(enum)

dedup(Enumerable.t()) :: Enumerable.t()

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

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

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

Примеры

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

dedup_by(enum, fun)

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)

drop(Enumerable.t(), non_neg_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)

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)

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]

each(enum, fun)

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

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

Полезно для добавления побочных эффектов (например, печати) в поток.

Примеры

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)

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)

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]]

intersperse(enumerable, intersperse_element)

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)

interval(non_neg_integer()) :: 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)

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

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

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

iterate(start_value, next_fun)

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

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

Примеры

iex> Stream.iterate(0, &(&1+1)) |> Enum.take(5)
[0, 1, 2, 3, 4]

map(enum, fun)

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)

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)

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)

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

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

Примеры

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

resource(start_fun, next_fun, after_fun)

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)

run(stream)

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

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

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

Примеры

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

stream = File.stream!("code")
|> Stream.map(&String.replace(&1, "#", "%"))
|> Stream.into(File.stream!("new"))
|> Stream.run

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

scan(enum, fun)

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)

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)

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)

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)

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)

timer(non_neg_integer()) :: Enumerable.t()

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

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

Примеры

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

transform(enum, acc, reducer)

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 = 1..100
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)
[1, 2, 3]

transform(enum, start_fun, reducer, after_fun)

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

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

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

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

unfold(next_acc, next_fun)

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]

uniq(enum)

uniq(Enumerable.t()) :: Enumerable.t()

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

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

Примеры

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

uniq_by(enum, fun)

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, offset \\ 0)

with_index(Enumerable.t(), integer()) :: Enumerable.t()

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

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

Примеры

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}]

zip(enumerables)

zip(Enumerable.t()) :: Enumerable.t()
zip([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(left, right)

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

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

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

Примеры

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}]

© 2012 Plataformatec
Licensed under the Apache License, Version 2.0.
https://hexdocs.pm/elixir/1.6.6/Stream.html

Spec-Zone.ru

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