【问题标题】:For Loop for list of Objects with use Multithread in PythonFor循环用于在Python中使用多线程的对象列表
【发布时间】:2019-12-12 16:39:16
【问题描述】:

我是 Python 新手,我有一个程序可以加载一个超过 100k 行的大 CSV 文件,每行有 4 列。 在 FOR 循环中,我检查每一行是否有相同的重复列表 (dlist),这个 dlist 是我与另一个加载的 DRef 类的对象列表功能

DsRef 类:

from tqdm import tqdm
from multiprocessing import Pool, cpu_count, freeze_support

class DsRef:
    def __init__(self, pn, comp, comp_name, type, diff):
        self.pn = pn
        self.comp = comp
        self.comp_name = comp_name
        self.type = type
        self.diff = diff

    def __str__(self):
        return f'{self.pn} {get_red("|")} {self.comp} {get_red("|")} {self.comp_name} {get_red("|")} {self.type} {get_red("|")} {self.diff}\n'

    def __repr__(self):
        return str(self)  

    def __iter__(self):
        return iter(self.__dict__.items())

复制类:

class Duplication:
    def __init__(self, pn, comp, cnt):
        self.pn = pn
        self.comp = comp
        self.cnt = cnt

    def __str__(self):
        return f'{self.pn};{self.comp};{self.cnt}\n'

    def __repr__(self):
        return str(self)

    def __hash__(self):
        return hash(('pn', self.pn,
                 'comp', self.comp))

    def __eq__(self, other):
        return self.pn == other.pn and self.comp == other.comp 

加载数据文件样本进行测试:

dlist= []
dlist.append(DsRef(
                    "TTT_XXX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XCX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XXX", "CCC_VCV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XXX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XYX", "CCC_YYY", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TAT_XQX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "ATT_XXX", "CCC_VQV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_EEE", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XWX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_XXX", "CCC_VWV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
                    "TTT_EEE", "CCC_VVV", "CTYPE", "CTYPE", "text"))

查找并返回重复值行的方法:

def FindDuplications(dlist):
    duplicates = []
    for pn, comp in enumerate(dlist):            
        matches = [xpn for xpn, xcomp in enumerate(dlist) if pn == xpn and comp == xcomp]
        duplicates.append(Duplication(pn, comp, len(matches)))
    return duplicates

row.pn == x.pn and row.comp == x.comp 如果它是真的我发现一个重复我将每个 objech 的前 2 个参数与列表中的每个对象进行比较

现在我尝试使用类似的东西来使用所有处理器以获得更快的结果,现在需要超过 15 分钟

if __name__ == '__main__':
    freeze_support()
    p = Pool(cpu_count())
    duplicates = p.map(FindDuplications, dlist)
    p.close()
    p.join()

首先,当 Class 不可迭代时我得到一个错误,然后我为第一个类创建 iter 函数,之后,我得到一个错误,然后元组对象不知道 pncomp 参数,然后我使用 in for enumerate(dlist) 但仍然不起作用

你能帮帮我吗?

我也想使用 TQDM 来检查处理功能的进度以查找重复


有一个不使用多线程的原始工作函数:

def CheckDuplications(dlist):
    print(get_yellow("========= CHECK CROSS DUPLICATIONS ========="))
    duplicates = []
    for r in tqdm(dlist):
        matches = [x for x in dlist if r.pn == x.pn and r.comp == x.comp]
        duplicates.append(Duplication(r.pn, r.comp, len(matches)))

    results = [d for d in duplicates if d.cnt > 1]
    results = set(results)
    return results

从函数 FindDuplications 我得到了 DsRef 对象的列表(简单副本),但这必须返回 Duplication 对象的列表,有问题

谢谢

【问题讨论】:

  • 到底是什么问题?我刚刚尝试运行您的代码,它似乎至少执行得很好。输出不是你所期望的吗?
  • 当我尝试运行 CheckDuplications 函数时,它工作正常,但它只使用 12 个逻辑处理器内核中的一个,当我使用 FindDuplications 时,我会喜欢使用所有逻辑核心或多个逻辑核心以获得更快的结果。当我运行多线程函数 FindDuplications 脚本在 3 秒后结束,但函数 CheckDuplications 需要超过 17 分钟,但我有超过 100k 行
  • 哦,我想我遇到了你的问题。 pool.map() 独立调用每个项目的给定函数。 FindDuplications 没有收到完整列表,也无法访问列表的其余部分以查找其他重复项。
  • 顺便说一句,python约定使用snake_case作为函数,应该是find_duplications
  • 好的,snake_case 会好的,但是您知道如何解决这个问题或如何解决它吗?

标签: python-3.x multithreading list


【解决方案1】:

代码中有一些问题,你没有并行它,你不能只在多核上运行繁重任务的单线程代码。代码需要一些采纳。

好吧,不管怎样,我们到了:)

from math import ceil
from multiprocessing import Pool, cpu_count, freeze_support


def get_red(val):
    return val


class DsRef:
    def __init__(self, pn, comp, comp_name, type, diff):
        self.pn = pn
        self.comp = comp
        self.comp_name = comp_name
        self.type = type
        self.diff = diff

    def __str__(self):
        return f'{self.pn} {get_red("|")} {self.comp} {get_red("|")} {self.comp_name} {get_red("|")} {self.type} {get_red("|")} {self.diff}\n'

    def __repr__(self):
        return str(self)


class Duplication:
    def __init__(self, pn, comp, cnt):
        self.pn = pn
        self.comp = comp
        self.cnt = cnt

    def __str__(self):
        return f'{self.pn};{self.comp};{self.cnt}\n'

    def __repr__(self):
        return str(self)

    def __hash__(self):
        return hash(('pn', self.pn,
                     'comp', self.comp))

    def __eq__(self, other):
        return self.pn == other.pn and self.comp == other.comp


dlist = []
dlist.append(DsRef(
    "TTT_XXX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XCX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XXX", "CCC_VCV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XXX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XYX", "CCC_YYY", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TAT_XQX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "ATT_XXX", "CCC_VQV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_EEE", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XWX", "CCC_VVV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_XXX", "CCC_VWV", "CTYPE", "CTYPE", "text"))
dlist.append(DsRef(
    "TTT_EEE", "CCC_VVV", "CTYPE", "CTYPE", "text"))


def FindDuplications(task):
    dlist, start, count = task

    duplicates = []
    for r in dlist[start:start + count]:
        matches = [x for x in dlist if r.pn == x.pn and r.comp == x.comp]
        duplicates.append(Duplication(r.pn, r.comp, len(matches)))

    return {d for d in duplicates if d.cnt > 1}


if __name__ == '__main__':
    freeze_support()

    threads = cpu_count()
    tasks_per_thread = ceil(len(dlist) / threads)

    tasks = [(dlist, tasks_per_thread * i, tasks_per_thread) for i in range(threads)]

    p = Pool(threads)
    duplicates = p.map(FindDuplications, tasks)
    p.close()
    p.join()

    duplicates = {item for sublist in duplicates for item in sublist}

    print(duplicates)
    print(type(duplicates))

它对我来说效果很好,返回的结果与单线程函数相同,并且可以在所有可用内核中并行工作。

输出

python test.py
{TTT_EEE;CCC_VVV;2
, TTT_XXX;CCC_VVV;2
}
<class 'set'>

【讨论】:

  • 嗨,Alexander,感谢它现在与所有内核并行工作,但我只需要返回 Duplication Class,而不是 DsRef。进入复制类,我需要加载所有数据,在我选择数据后,类属性 CNT 大于 1 以仅选择复制。这意味着,然后每个项目表单任务列表(DsRef 类对象)必须检查 dlist 中的所有行。是否可以使用 VS Code 断点进行并行调试?
  • 源代码适用于您的函数,只是针对多线程使用进行了修改。我没有更改与逻辑相关的任何内容duplicates.append(Duplication(pn, comp, len(matches)))
  • 原始(非多线程)函数返回 Duplication 对象列表,但现在返回 Duplicate 对象,其中属性 comp 具有 DsRef 类,但 Duplication 类只有字符串属性和重复次数的 CNT 属性跨度>
  • woof,我没看到你改变了原来的功能。 OK,更新到原来的逻辑了。请立即查看
  • 我什么都没有改变,但没关系。我需要在函数中添加 2 个参数(首先是每个核心的任务列表,也是完整的 dlist )因为任务列表必须检查原始列表中的所有数据
猜你喜欢
  • 2015-06-23
  • 2011-02-13
  • 2013-04-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-23
  • 1970-01-01
相关资源
最近更新 更多