【问题标题】:How is the PLINQ AsParallel function passing data to a function when called within the same scope在同一范围内调用时,PLINQ AsParallel 函数如何将数据传递给函数
【发布时间】:2019-10-30 13:11:26
【问题描述】:

使用 IronPython,我正在并行调用一些函数,该函数在并行化数据所在的同一函数内,以使其保持在同一范围内。

在 CPython 的多处理中,数据必须显式传递给子进程、打开多少子进程等是非常清楚的。这很容易理解开销。

对于 PLINQ,代码如何并行运行?即:

是否有另一个 Ironpython 实例正在运行并且所有内容都再次导入?例如 import myHugeLibrary 将在每次创建文件的新 python 实例时运行。

CalcParallel() 接收一些数据数组和一个字典。在这个范围内是一个应该并行运行的函数computation(),它在主脚本中调用另一个函数checkVals()。由于computation() 与调用AsParallel() 处于同一范围内,因此我不需要显式地将要使用的数据传递给它。然而,这是否意味着数据被复制到每个进程/线程,或者作为参考保存,并且在仅被读取(而不是写入)时很好?如果它被复制,是否每次计算一个项目时都复制,这意味着如果列表中有 100 个项目和 10 个线程,它将复制数据 10 次,因为它将 100 个项目放入 10 个块中,还是复制 100 次?

同样,示例 C_dict 数据在计算了一些数据之后,在运行下一轮数据之前进行了修改(根据结果,它添加了更多的 todo)。当并行进程运行时,是否会再次复制修改后的数据?

下面是我想知道的一些示例代码结构。它并不是关于代码本身,但我写这个只是为了说明问题,即使它不是正确的方法。

# get LINQ dependencies
import clr
clr.AddReference("System.Core")
import System
clr.ImportExtensions(System.Linq)
from System.Threading.Tasks import *

#import some huge library that takes time
import myHugeLibrary

max_val = 4 #some global value used within the thread

def checkVals(itemToCheck,A_vals,B_vals):
    #check against some global value
    if itemToCheck < max_val:
        return 0
    #do something else with A_vals

def CalcParallel(todo_list,A_vals,B_vals,C_dict): 
    """
    take in some data that is used in the functions that will
    run in parallel.
    """

    total_list = []

    #make a function that will be run in parallel
    def computation(itemToCheck):
        checkedItems = checkVals(itemToCheck,A_vals,B_vals)
        results = []
        for item in checkedItems: 
                results.append(item)
        return results

    #in a loop keep sending something out for calculation in parallel until
    # all the combinations are done
    while len(todo_list) != 0:
            #use AsParallel on a list of items
            results = todo_list.AsParallel().SelectMany(
                            lambda itemToCheck: 
                                computation(itemToCheck) ).ToList()

            todo_list = []
            for item in results:
                if item not in total_list: 
                    total_list.append(item)

                    #do some modification to the dictionary that was passed in
                    C_dict[item] = None

    return total_list


def main():
    todo_list = [3,3,2,4,5,4,1,3,4,5,1]
    A_vals = [0,1,2,3,4,5,6]
    B_vals = [-1,-3,-5,-7,-9]
    C_dict = {0:-3,4:-7}

    newVals = CalcParallel(todo_list,A_vals,B_vals,C_dict)

    print(newVals)

main()

【问题讨论】:

    标签: .net ironpython plinq


    【解决方案1】:

    PLINQ 的行为有点类似于 TPL。使用相对轻量级的结构对工作进行分区/计划/分派。没有额外的 Ironpython 进程,并且工作很可能安排在工作线程池上(通过任务),这意味着开销应该非常小。

    您使用的所有内容都由作用域/引用捕获,您应该通过使用线程安全的共享集合或者更好的是,通过以最终可以合并结果的方式对数据流进行建模来避免写入操作发生冲突(就像您对 SelectMany 所做的那样)。这不应导致您的数据被复制。

    为了确保高效执行,工作块的大小需要合理以避免意外开销。

    【讨论】:

    • 您好,再次感谢您的信息!如果因为在范围内而没有复制任何内容,那么 .tolist() 函数只是将结果聚合到一个列表中,对吗?那么开销涉及到什么,仅仅是线程的启动还是还有其他方面的呢?
    猜你喜欢
    • 1970-01-01
    • 2017-06-24
    • 1970-01-01
    • 1970-01-01
    • 2011-06-30
    • 2014-12-05
    • 1970-01-01
    • 1970-01-01
    • 2019-05-12
    相关资源
    最近更新 更多