【发布时间】: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