使用@sync @spawn多线程批量处理时+=操作符异常及任务重复问题
Julia多线程批量处理重复执行导致计算错误的问题
问题现象
在手动实现多线程批量处理数组时,使用+=操作符出现结果异常。单线程运行测试脚本结果正确,但以3线程模式(执行命令:julia --threads 3 ./parallel_issues_test.jl)运行时,部分索引的计算结果不符合预期,错误输出如下:
Error at index CartesianIndex(1, 6): 10.5 != 7 Error at index CartesianIndex(1, 11): 18.0 != 12
测试代码:
using Base.Threads function function_1(cartesian_index) i = cartesian_index[1] j = cartesian_index[2] return (i + j)/2 end function main() array_1 = zeros(Float64, 15, 15) array_2 = zeros(Float64, 15, 15) array_3 = zeros(Float64, 15, 15) cartesian_indecies = CartesianIndices(array_1) number_of_indecies = length(cartesian_indecies) n_threads = Threads.nthreads() batch_size = ceil(Int, number_of_indecies / n_threads) cartesian_indicies = CartesianIndices(array_1) @sync for batch_index in 1:batch_size:number_of_indecies Threads.@spawn begin for view_index in batch_index:min(number_of_indecies, batch_index + batch_size) cartesian_index = cartesian_indecies[view_index] array_1[cartesian_index] = function_1(cartesian_index) array_2[cartesian_index] = function_1(cartesian_index) end end end for index in CartesianIndices(array_1) if array_1[index] != (index[1] + index[2])/2 println("Error part 1 at index $index: $(array_1[index]) != $(index[1] + index[2])") end end @sync for batch_index in 1:batch_size:number_of_indecies Threads.@spawn begin for view_index in batch_index:min(number_of_indecies, batch_index + batch_size) cartesian_index = cartesian_indecies[view_index] array_1[cartesian_index] += function_1(cartesian_index) array_3[cartesian_index] = function_1(cartesian_index) end end end for index in CartesianIndices(array_1) if array_1[index] != index[1] + index[2] println("Error at index $index: $(array_1[index]) != $(index[1] + index[2])") end end array_4 = array_2 + array_3 for index in CartesianIndices(array_4) if array_4[index] != index[1] + index[2] println("Array 4 Error at index $index: $(array_1[index]) != $(index[1] + index[2])") end end end main()
问题根源
手动划分batch的逻辑存在错误:Julia的范围是左闭右闭区间,代码中使用batch_index:min(number_of_indecies, batch_index + batch_size)作为每个batch的索引范围,导致相邻batch的索引重叠。
以batch_size=75(15×15数组共225个索引,3线程下ceil(225/3)=75)为例:
- 第一个batch的范围是
1:76(1+75=76) - 第二个batch的起始索引为76,范围是
76:151
索引76会被两个线程重复处理,对array_1执行两次+= function_1(cartesian_index)操作,最终结果变为(i+j)/2 + 2*(i+j)/2 = 3*(i+j)/2,与预期的i+j不符。
修复方案
方案1:修正batch索引范围
将每个batch的结束索引改为batch_index + batch_size - 1,确保索引不重叠:
修改后的并行循环代码:
# 第一阶段赋值 @sync for batch_index in 1:batch_size:number_of_indecies Threads.@spawn begin end_index = min(number_of_indecies, batch_index + batch_size - 1) for view_index in batch_index:end_index cartesian_index = cartesian_indecies[view_index] array_1[cartesian_index] = function_1(cartesian_index) array_2[cartesian_index] = function_1(cartesian_index) end end end # 第二阶段累加 @sync for batch_index in 1:batch_size:number_of_indecies Threads.@spawn begin end_index = min(number_of_indecies, batch_index + batch_size - 1) for view_index in batch_index:end_index cartesian_index = cartesian_indecies[view_index] array_1[cartesian_index] += function_1(cartesian_index) array_3[cartesian_index] = function_1(cartesian_index) end end end
另外,代码中cartesian_indecies存在拼写错误(正确应为cartesian_indices),建议修正以提升可读性。
方案2:使用内置并行宏简化实现
Julia提供了Threads.@threads宏,可自动分配线程任务,无需手动划分batch,代码更简洁且不易出错:
替换原手动分batch的代码为:
# 第一阶段赋值 Threads.@threads for cartesian_index in CartesianIndices(array_1) array_1[cartesian_index] = function_1(cartesian_index) array_2[cartesian_index] = function_1(cartesian_index) end # 第二阶段累加 Threads.@threads for cartesian_index in CartesianIndices(array_1) array_1[cartesian_index] += function_1(cartesian_index) array_3[cartesian_index] = function_1(cartesian_index) end
内容的提问来源于stack exchange,提问作者Jhayes2118
相关产品推荐
相关产品推荐

