【问题标题】:multiprocessing.Queue() into dictionary with large input size将 multiprocessing.Queue() 放入具有大输入大小的字典中
【发布时间】:2022-01-23 21:50:56
【问题描述】:

我有一个名为 encrypted_messages 的列表,其中包含 126018 个字符串。每个字符串都是加密的消息。我还有一个名为 decipher 的函数,给定一个字符串和一个密钥(包括 9 到 15 的整数),返回解密的消息。我需要使用每个密钥解密每条消息。由于 decipher 函数的计算量很大并且消息很多,因此我实现了一个多处理解决方案。 我创建了一个名为 messages_queue 的 multiprocessing.JoinableQueue() 包含所有加密的消息,并创建了一个名为 results_queue 的 multiprocessing.Queue() 来存储结果。这些队列由所有进程共享。进程从 messages_queue 获取消息,使用所有密钥对其应用 decipher 并将结果存储为 2 个元素的列表(用于解密消息的密钥和解密的消息)。它看起来像这样:

[9, message_1], [15, message_2], [14, message_3], ...

results_queue 有 882126 个元素,正如预期的那样(注意 126018*7 = 882126),其中每个元素都是一个列表。 我想从 results_queue 中获取长度为 7 的字典,其中每个键都是一个整数,每个值都是一个列表,其中包含使用该键解密的所有消息。它应该是这样的:

{9:[decrypted messages using key 9], 10:[decrypted messages using key 10], ...,
15:[decrypted messages using key 15]}

我已经尝试了几种方法来做到这一点,但我无法提出解决方案。我分享以下代码:

final_results = {key:[] for key in range(9, 16)}
while not results_queue.empty():
    message = results_queue.get() # Note that this is a list: [key, message]
    final_results[message[0]].append(message[1])

我也尝试过先创建一个这样的列表(我可以从列表中创建字典):

results = []
results_queue.put('STOP')
while True:
    message = results_queue.get()
    if message == 'STOP':
        break
    results.append(message)

我也尝试过像这样使用带有哨兵的迭代器:

results = []
results_queue.put(None)
for message in iter(results_queue.get, None):
    results.append(message)

通过所有这些方法,我丢失了很多(超过 50%)的消息。该列表应该有 882126 个列表,并且每次我运行代码时它都有一个不同且更小的数字。这个数字在我看来完全是随机的。我不知道如何解决这个问题,因为当我使用更小的列表(例如 100 个元素)时,上述方法可以正常工作。 这个问题与输入大小有关吗?我的 multiprocessing.Queue() 是否太大?我认为这不是进程之间的协调问题,因为我获得的 Queue() 是我所期望的,然后进程结束,但也许我遗漏了一些东西。

如果有用的话,我使用的是 Python 3.8.5 和 Linux Mint 20.2。欢迎任何帮助,因为我有点卡住了。提前致谢。

【问题讨论】:

  • 进程同步可能存在一些问题。您是否考虑过在填充 results_queue 的进程中添加终止元素?可能您在仍然生成元素时放置 None/'STOP',因此它不在队列的末尾。
  • 可能是这样,但看起来很奇怪,因为 results_queue() 是由函数返回的。所以,只有当函数完成执行时,我才能得到队列。它具有正确的结构,因为 results_queue.qsize() 给了我 882126,当我运行 results_queue.get() 时,我得到了预期的列表:[9,message_decrypted using key 9],等等。它是在调用函数之后,当返回队列,我将 None/'STOP' 放入其中。不管怎样,我回家后会尝试你的建议。谢谢你的回答。

标签: python dictionary python-multiprocessing


【解决方案1】:

这是创建具有这样的 for 的字典的代码

{key1:[message_1,message_2],key2:[message_3,message_4]}

message_decoded 的形状必须是

[[key1,message1],[key2,message2]]

dict = {}

messages_decoded = []

for item in messages_decoded:

    if item[0] in dict:
        dict[item[0]].append(item[1])
    else:
        dict[item[0]] = [item[1]]

编辑

此代码将结果队列转换为列表。

list_messages = [results_queue.get() for _ in range(results_queue.qsize())]

【讨论】:

  • 您好,感谢您的回答。我想我的问题是获得具有您指出的形状的 messages_decoded 。我拥有的是一个 multiprocessing.Queue(),而不是一个列表。我尝试使用 results_queue.get() 方法遍历它并将这些元素附加到列表中,但是很多元素都丢失了(未附加)。正如我所说,只有当元素很多时才会出现这个问题。
  • @jakeis 我将队列的转换添加到我的代码中的列表中
  • 它会引发错误。 type(results_queue) = 。如果我尝试做 list(results_queue) 我得到以下信息: TypeError: 'Queue' object is not iterable
  • @jakeis 我创建了一个应该可以工作的新方法 :) 我用 10**5 个元素对其进行了测试,它就像一个魅力 :)
  • 确实,它就像一个魅力。这是一个聪明的解决方案。非常感谢!
猜你喜欢
  • 1970-01-01
  • 2012-04-19
  • 1970-01-01
  • 2019-05-31
  • 2019-05-10
  • 2016-11-15
  • 1970-01-01
  • 2016-01-01
  • 1970-01-01
相关资源
最近更新 更多