【问题标题】:Task.async in Elixir StreamElixir Stream 中的 Task.async
【发布时间】:2015-12-11 21:30:09
【问题描述】:

我想在一个大列表上做一个平行映射。代码看起来有点像这样:

big_list
|> Stream.map(&Task.async(Module, :do_something, [&1]))
|> Stream.map(&Task.await(&1))
|> Enum.filter filter_fun

但我正在检查 Stream 实现,据我了解 Stream.map 组合了函数并将组合函数应用于流中的元素,这意味着序列是这样的:

  1. 取第一个元素
  2. 创建异步任务
  3. 等待它完成
  4. 拿第二个元素...

在这种情况下,它不会并行执行。我是对的还是我错过了什么?

如果我是对的,那么这段代码呢?

Stream.map Task.async ...
|> Enum.map Task.await ...

这会并行运行吗?

【问题讨论】:

标签: parallel-processing elixir


【解决方案1】:

第二个也没有做你想做的事。用这段代码可以看得很清楚:

defmodule Test do
  def test do
    [1,2,3]
    |> Stream.map(&Task.async(Test, :job, [&1]))
    |> Enum.map(&Task.await(&1))
  end

  def job(number) do
    :timer.sleep 1000
    IO.inspect(number)
  end
end

Test.test

您会看到一个数字,然后等待 1 秒,然后再看到一个数字,以此类推。这里的关键是你想尽快创建任务,所以你不应该使用 懒惰Stream.map。取而代之的是使用急切的Enum.map

|> Enum.map(&Task.async(Test, :job, [&1]))
|> Enum.map(&Task.await(&1))

另一方面,您可以在等待时使用Stream.map,只要您稍后进行一些急切的操作,例如您的filter。这样,等待将穿插在您可能对结果进行的任何处理中。

【讨论】:

    【解决方案2】:

    Elixir 1.4 提供了新的 Task.async_stream/5 函数,该函数将返回一个流,该流在可枚举中的每个项目上同时运行给定函数。

    还可以使用:max_concurrency:timeout 选项参数指定最大工作线程数和超时。

    请注意,您不必等待此任务,因为该函数返回一个流,因此您可以使用Enum.to_list/1 或使用Stream.run/1


    这将使您的示例同时运行:

    big_list
    |> Task.async_stream(Module, :do_something, [])
    |> Enum.filter(filter_fun)
    

    【讨论】:

      【解决方案3】:

      你可以试试Parallel Stream

      stream = 1..10 |> ParallelStream.map(fn i -> i * 2 end)
      stream |> Enum.into([])
      [2,4,6,8,10,12,14,16,18,20]
      

      UPD 或者更好地使用Flow

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2020-04-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-08-22
        • 1970-01-01
        • 1970-01-01
        • 2018-09-04
        相关资源
        最近更新 更多