Julia 并行计算

FreeGuideOnline 最新 2026-07-13

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) 会等待任务完成并获取返回值。这种模式非常灵活,适合实现分治算法。

分布式并行:@distributedpmap

分布式并行基于多进程,每个进程拥有独立的内存空间,数据通过消息传递共享。核心模块是 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),甚至可以自定义归约函数。

手动管理任务:remotecallfetch

对于更精细的控制,你可以直接向指定 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.jlMetal.jl (Apple)。以 CUDA 为例:

using CUDA
# 在 GPU 上创建数组
a = CUDA.rand(1000, 1000)
b = CUDA.rand(1000, 1000)
c = a * b   # 自动在 GPU 上执行矩阵乘法