【问题标题】:Julia: parallelize operations on complex data structures (eg DataFrames)Julia:对复杂数据结构(例如 DataFrames)进行并行操作
【发布时间】:2020-09-08 09:06:11
【问题描述】:

我想并行处理多个大型数据集。不幸的是,我从使用 Threads.@threads 获得的加速是非常次线性的,如下面的简化示例所示。

(我对 Julia 很陌生,如果我遗漏了一些明显的东西,请道歉)

让我们创建一些虚拟输入数据 - 8 个数据帧,每个数据帧有 2 个整数列和 1000 万行:

using DataFrames

n = 8
dfs = Vector{DataFrame}(undef, n)
for i = 1:n
    dfs[i] = DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7))))
end

现在对每个数据帧进行一些处理(按 x1 和 sum x2 分组)

function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

最后,比较在单个数据帧上执行处理的速度与在所有 8 个数据帧上并行执行处理的速度。我正在运行它的机器有 50 个内核,而 Julia 启动时有 50 个线程,所以理想情况下不应该有太大的时间差异。

julia> dfs_res = Vector{DataFrame}(undef, n)

julia> @time for i = 1:1
           dfs_res[i] = process(dfs[i])
       end
  3.041048 seconds (57.24 M allocations: 1.979 GiB, 4.20% gc time)

julia> Threads.nthreads()
50

julia> @time Threads.@threads for i = 1:n
           dfs_res[i] = process(dfs[i])
       end
  5.603539 seconds (455.14 M allocations: 15.700 GiB, 39.11% gc time)

因此,每个数据集的并行运行时间几乎是两倍(数据集越多,情况就越糟)。我感觉这与低效的内存管理有关。第二次运行的 GC 时间相当长。而且我假设undef 的预分配对于DataFrames 无效。我在 Julia 中看到的几乎所有并行处理示例都是在具有固定和先验已知大小的数字数组上完成的。然而,这里的数据集可以有任意大小、列等。在 R 工作流中,可以使用mclapply 非常有效地完成。 Julia 中是否有类似的(或不同但有效的模式)?我选择使用线程而不是多处理以避免复制数据(Julia 似乎不支持 R / mclapply 之类的 fork 进程模型)。

【问题讨论】:

  • 你设置JULIA_NUM_THREADS了吗? Threads.nthreads() 为你输出了什么?
  • 是的,我做了,见上面的输出。这是50
  • 请注意,在您的示例中,如果您使用:x2 => sum 而不是[:x2] => sum,则操作要快得多并且分配要少得多。这是因为 DataFrames 有一个快速的路径,可以在单个列上进行常见的缩减。我们可能会改进它以覆盖[:x2] => sum,但如果您知道自己只有一列,通常最好不要使用向量。 (另外,您的示例非常极端,因为每个组平均只有一行。所以它并不能真正代表我想说的大多数工作流——更多的组意味着更多的分配,所以更长的 GC 时间,这不是多线程的.)

标签: julia


【解决方案1】:

Julia 中的多线程不能很好地扩展到 16 线程之外。 因此,您需要改用多处理。 您的代码可能如下所示:

using DataFrames, Distributed
addprocs(4) # or 50
@everywhere using DataFrames, Distributed

n = 8
dfs = Vector{DataFrame}(undef, n)
for i = 1:n
    dfs[i] = DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7))))
end

@everywhere function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

dfs_res = @distributed (vcat) for i = 1:n
      df = process(dfs[i])
      (i, myid(), df)
end

在这种类型的代码中重要的是在进程之间传输数据需要时间。因此,有时您可能只想将 DataFrames 放在单独的工作人员身上。像往常一样 - 这取决于您的处理架构。

编辑一些关于性能的注释

为了测试,将您的代码放入函数中并使用consts(或使用 BenchamrTools.jl)

using DataFrames

const dfs = [DataFrame(Dict("x1" => rand(1:Int64(1e7), Int64(1e7)), "x2" => rand(1:Int64(1e7), Int64(1e7)))) for i in 1:8 ]

function process(df::DataFrame)::DataFrame
    combine([:x2] => sum, groupby(df, :x1))
end

function p1!(res, d)
    for i = 1:8
        res[i] = process(dfs[i])
    end
end


function p2!(res, d)
     Threads.@threads for i = 1:8
        res[i] = process(dfs[i])
    end
end

const dres = Vector{DataFrame}(undef, 8)

这里是结果

julia> GC.gc();@time p1!(dres, dfs)
 30.840718 seconds (507.28 M allocations: 16.532 GiB, 6.42% gc time)

julia> GC.gc();@time p1!(dres, dfs)
 30.827676 seconds (505.66 M allocations: 16.451 GiB, 7.91% gc time)

julia> GC.gc();@time p2!(dres, dfs)
 18.002533 seconds (505.77 M allocations: 16.457 GiB, 23.69% gc time)

julia> GC.gc();@time p2!(dres, dfs)
 17.675169 seconds (505.66 M allocations: 16.451 GiB, 23.64% gc time)

为什么在 8 核机器上差异只有大约 2 倍 - 因为我们大部分时间都在垃圾收集上! (查看您问题中的输出 - 问题是一样的) 当您使用更少的 RAM 时,您会看到更好的多线程加速,最高可达 3 倍。

【讨论】:

  • 谢谢!尽管我希望避免多处理的复制开销(我的实际用例要大得多),但对于此示例,这种方法似乎确实更快。另请注意,在我上面的示例中,有效线程数为 8,因此远低于您所说的 16 是不可行的。
  • 在您的示例中,您错误地测量了时间,因为您同时测量了编译时间和运行时。 Julia 的编译时间对于多线程代码比单线程代码要长得多,而对于分布式代码编译则需要更长的时间。看看:stackoverflow.com/questions/60867667/… 无论哪种情况,当您开始正确测量时,您都会对高达 16 个线程的性能感到满意
  • 我确实执行了每个命令几次 - 这是否意味着它会在第一次之后使用缓存的编译版本?否则我怎么能更好地衡量它?
  • 我添加了一些关于性能测试的 cmets。无论如何在 50 个内核上运行,你应该使用分布式而不是线程,有足够的内存以避免过多的 GC,并且在工作人员之间有一个良好的工作流和数据分布策略。
  • 感谢有关更精确测量执行时间的指针。然而,基本的结论仍然是一样的——多线程的一个非常次线性的缩放。我认为这与内核/线程的数量没有任何关系,因为这是一个故意使用仅 8 个的小示例。同意 GC。那么问题来了——如何避免呢?在此示例中,我仅使用 DataFrames 库的基本功能。
猜你喜欢
  • 2020-06-24
  • 1970-01-01
  • 1970-01-01
  • 2012-07-17
  • 1970-01-01
  • 2022-11-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多