Julia 并行计算
bash julia -t 4
或者在 Windows 下:
```powershell
julia --threads 4
你可以在 Julia 环境中检查当前可用的线程数:
Threads.nthreads()
如果你想在脚本或 Jupyter 中动态调整线程数,可以在代码开头调用(仅限未启动多线程时设置环境变量):
ENV["JULIA_NUM_THREADS"] = 4
# 但这只能在 Julia 会话启动前生效,更推荐在启动时指定
启动分布式 Worker 进程
分布式并行需要多个 Julia 进程。你可以通过以下方式添加 worker:
using Distributed
addprocs(4) # 添加 4 个本地 worker 进程
nprocs() # 返回总进程数(包括主进程)
nworkers() # 返回 worker 数量
也可以在启动时直接添加:
julia -p 4
这两种方式的效果相同,都会在后台启动额外的 Julia 进程。如果你需要跨机器集群,可以使用 --machinefile 指定机器列表,但本节我们聚焦本地多进程。
多线程并行:Threads.@threads
Julia 的多线程模型基于共享内存,所有线程共享变量,因此写并行代码时要特别注意数据竞争(race conditions)。
静态循环并行化 Threads.@threads
最直接的方式是用 Threads.@threads 宏并行化一个 for 循环。它会将循环迭代静态地分配给各个线程。例如,计算一个数组中每个元素的平方:
using Base.Threads
function parallel_square!(A)
@threads for i in eachindex(A)
A[i] = A[i]^2
end
return A
end
A = rand(1:10, 10^6)
parallel_square!(A)
这里的 @threads 要求循环体是独立的(无迭代间依赖),否则会产生不可预测的结果。
使用原子操作保证安全
当多个线程需要更新同一个变量时,必须使用原子操作来避免竞争。例如,多个线程累加同一个标量:
acc = Atomic{Int64}(0)
@threads for i in 1:1000
atomic_add!(acc, i)
end
println(acc[]) # 输出 500500
Julia 提供了 Atomic 类型和一系列原子操作(atomic_add!, atomic_cas! 等)。共享数组元素则需要使用 Threads.Atomic 包装。
动态任务调度:Threads.@spawn
对于不规则的并行任务(例如递归或任务数量多于线程数),可以使用 @spawn 动态创建任务。这些任务由 Julia 的运行时调度到可用线程上执行:
function fib(n)
n < 2 && return n
t = @spawn fib(n-1)
return fib(n-2) + fetch(t)
end
@time fib(30)
fetch(t) 会等待任务完成并获取返回值。这种模式非常灵活,适合实现分治算法。
分布式并行:@distributed 与 pmap
分布式并行基于多进程,每个进程拥有独立的内存空间,数据通过消息传递共享。核心模块是 Distributed。
基础概念:主进程 vs Worker 进程
- 主进程(process 1)负责协调,也可参与计算。
- Worker 进程运行在后台,执行分配的任务。
- 变量和函数需要通过
@everywhere在所有进程中定义,否则 worker 无法访问。
例如:
@everywhere using Statistics
@everywhere function myfunc(x)
return mean(x) .+ rand()
end
并行 map 操作:pmap
pmap 是最简单的分布式并行函数,类似于 Python 的 multiprocessing.Pool.map。它会对集合中的每个元素应用一个函数,并在不同 worker 上调度计算:
# 示例:对一个大向量的每个元素进行耗时计算
inputs = 1:100
@everywhere function slow_compute(x)
sleep(0.1)
return log(x + 1)
end
results = pmap(slow_compute, inputs)
pmap 会自动平衡负载,并收集结果返回给主进程。很适合任务执行时间相近的场景。
并行归约:@distributed
@distributed 宏用于并行化归约操作(例如求和、求最大值)。它可以将一个迭代空间划分到多个进程,然后合并结果。
语法:
@distributed (归约运算符) for i in range
# 计算单个元素
end
例如,计算 π 的蒙特卡洛近似:
@everywhere function monte_carlo_pi(N)
inside = 0
for _ in 1:N
x = rand()
y = rand()
inside += (x^2 + y^2 <= 1)
end
return 4 * inside / N
end
total_samples = 10^8
per_worker = div(total_samples, nworkers())
# 使用 @distributed 将 per_worker 次试验分配到各 worker,最后用 + 合并结果
pi_est = @distributed (+) for i in 1:nworkers()
monte_carlo_pi(per_worker)
end
println(pi_est)
(+) 表示各 worker 的结果按 + 归约。你也可以使用其他可结合的操作符,如 (*), (min), (max),甚至可以自定义归约函数。
手动管理任务:remotecall 与 fetch
对于更精细的控制,你可以直接向指定 worker 发送函数调用:
r = remotecall(myfunc, 2, 42) # 在 worker 2 上调用 myfunc(42)
result = fetch(r) # 等待并获取结果
或使用 @spawnat 宏:
f = @spawnat :any myfunc(100)
result = fetch(f)
:any 表示调度到任意空闲 worker。
共享数组:跨进程的高效数据共享
在分布式环境中,频繁传递大数组会带来显著通信开销。SharedArray 提供了一种将数组存储在共享内存段中,让所有 worker 都能直接读写的机制(仅限于单台机器)。
using Distributed, SharedArrays
addprocs(4)
# 创建共享数组,维度 (1000,),类型 Float64,初始化由函数完成
S = SharedArray{Float64}((1000,); init = S -> S[localindices(S)] = rand(length(localindices(S))))
每个 worker 通过 localindices(S) 获得自己负责的那部分索引(根据线程或进程自动划分)。这样每个 worker 只写入自己负责的区间,避免冲突。
如果需要进行并行读取和全局操作,可以结合 @distributed:
total = @distributed (+) for i in 1:length(S)
S[i]^2
end
注意:写入共享数组时,如果多个 worker 同时写同一内存位置,会导致数据竞争,需要自行协调或使用原子操作(Julia 1.7+ 增加 SharedArray 的原子支持,但一般通过分区避免竞争更简单)。
实际案例:并行图片处理
假设我们要对一组图片进行灰度转换并缩放。使用 Images.jl 加载图片,然后利用 pmap 并行处理。
using Distributed, Images, FileIO
@everywhere using Images, FileIO
# 假设 images/ 目录下有一系列 jpg 文件
files = filter(endswith(".jpg"), readdir("images/"; join=true))
resize_size = (128, 128)
@everywhere function process_image(path, size)
img = load(path)
gray_img = Gray.(img)
small_img = imresize(gray_img, size)
return small_img
end
processed = pmap(f -> process_image(f, resize_size), files)
# 保存结果
for (i, img) in enumerate(processed)
save("output/frame_$i.jpg", img)
end
这个模式可以直接扩展到视频帧处理、批量科学数据提取等场景。
GPU 并行简介
Julia 的 GPU 计算主要通过包实现,如 CUDA.jl (NVIDIA)、AMDGPU.jl、Metal.jl (Apple)。以 CUDA 为例:
using CUDA
# 在 GPU 上创建数组
a = CUDA.rand(1000, 1000)
b = CUDA.rand(1000, 1000)
c = a * b # 自动在 GPU 上执行矩阵乘法