【问题标题】:Troubleshooting data inconsistencies with Python multiprocessing/threading使用 Python 多处理/线程解决数据不一致问题
【发布时间】:2015-01-21 14:38:58
【问题描述】:

TL;DR:使用线程和多处理以及单线程运行代码后得到不同的结果。需要有关故障排除的指导。

您好,如果这可能有点过于笼统,我提前道歉,但我需要一些帮助来解决问题,我不确定如何最好地继续。

故事是这样的;我有一堆数据被索引到 Solr 集合(约 2.5 亿个项目)中,该集合中的所有项目都有一个 sessionid。某些项目可以共享相同的会话 ID。我正在梳理集合以提取具有相同会话的所有项目,稍微处理数据并吐出另一个 JSON 文件以供稍后索引。

代码有两个主要功能: proc_day - 接受一天并处理当天的所有会话 和 proc_session - 完成单个会话需要发生的所有事情。

多处理是在proc_day上实现的,所以每一天都会被一个单独的进程处理,proc_session函数可以用线程来运行。下面是我在下面用于线程/多处理的代码。它接受一个函数、一个参数列表和线程/多进程的数量。然后它将根据输入参数创建一个队列,然后创建进程/线程并让它们通过它。我没有发布实际代码,因为它通常运行良好的单线程没有任何问题,但如果需要可以发布它。

autoprocs.py

import sys
import logging
from multiprocessing import Process, Queue,JoinableQueue
import time
import multiprocessing
import os

def proc_proc(func,data,threads,delay=10):
    if threads < 0:
        return
    q = JoinableQueue()
    procs = []

    for i in range(threads):
        thread = Process(target=proc_exec,args=(func,q))
        thread.daemon = True;
        thread.start()
        procs.append(thread)

    for item in data:
        q.put(item)

    logging.debug(str(os.getpid()) + ' *** Processes started and data loaded into queue waiting')

    s = q.qsize()
    while s > 0:
        logging.info(str(os.getpid()) + " - Proc Queue Size is:" + str(s))
        s = q.qsize()
        time.sleep(delay)

    for p in procs:
        logging.debug(str(os.getpid()) + " - Joining Process {}".format(p))
        p.join(1)

    logging.debug(str(os.getpid()) + ' - *** Main Proc waiting')
    q.join()
    logging.debug(str(os.getpid()) + ' - *** Done')

def proc_exec(func,q):
    p = multiprocessing.current_process()
    logging.debug(str(os.getpid()) + ' - Starting:{},{}'.format(p.name, p.pid))
    while True:
        d = q.get()
        try:
            logging.debug(str(os.getpid()) + " - Starting to Process {}".format(d))
            func(d)
            sys.stdout.flush()
            logging.debug(str(os.getpid()) + " - Marking Task as Done")
            q.task_done()
        except:
            logging.error(str(os.getpid()) + " - Exception in subprocess execution")
            logging.error(sys.exc_info()[0])
    logging.debug(str(os.getpid()) + 'Ending:{},{}'.format(p.name, p.pid))

autothreads.py:

import threading
import logging
import time
from queue import Queue

def thread_proc(func,data,threads):
    if threads < 0:
        return "Thead Count not specified"
    q = Queue()

    for i in range(threads):
        thread = threading.Thread(target=thread_exec,args=(func,q))
        thread.daemon = True
        thread.start()

    for item in data:
        q.put(item)

    logging.debug('*** Main thread waiting')
    s = q.qsize()
    while s > 0:
        logging.debug("Queue Size is:" + str(s))
        s = q.qsize()
        time.sleep(1)
    logging.debug('*** Main thread waiting')
    q.join()
    logging.debug('*** Done')

def thread_exec(func,q):
    while True:
        d = q.get()
        #logging.debug("Working...")
        try:
            func(d)
        except:
            pass
        q.task_done()

在 python 在不同的多处理/线程配置下运行后,我遇到了验证数据的问题。有很多数据,所以我真的需要让多处理工作。这是我昨天的测试结果。

Only with multiprocessing - 10 procs: 
Days Processed  30
Sessions Found  3,507,475 
Sessions Processed 3,514,496 
Files 162,140 
Data Output: 1.9G

multiprocessing and multithreading - 10 procs 10 threads
Days Processed  30
Sessions Found   3,356,362 
Sessions Processed   3,272,402 
Files    424,005 
Data Output: 2.2GB

just threading - 10 threads
Days Processed  31
Sessions Found   3,595,263 
Sessions Processed   3,595,263 
Files    733,664 
Data Output: 3.3GB

Single process/ no threading
Days Processed  31
Sessions Found   3,595,263 
Sessions Processed   3,595,263 
Files    162,190 
Data Output: 1.9GB

这些计数是通过日志文件中的 grepping 和 counties 条目收集的(每个主进程 1 个)。跳出来的第一件事是处理的天数不匹配。但是,我手动检查了日志文件,似乎缺少一个日志条目,日志条目后面有表明这一天已实际处理。我不知道为什么它被省略了。

我真的不想写更多的代码来验证这个代码,这似乎是一种可怕的浪费时间,有没有其他选择?

【问题讨论】:

  • 我认为人们会提供帮助,但只是关于以下内容的评论:“我真的不想编写更多代码来验证此代码,这似乎是在浪费时间” i>:如果数据分析很重要,那么您最好 100% 确保您的处理管道按照预期的方式工作。相应的测试代码可能很容易比您正在测试的代码更复杂。
  • "Sessions Found 3,595,263 , Sessions Processed 3,595,263 " -- 这是预期的输出吗?意思是,“单进程/无线程”和“仅线程 - 10 个线程”是否按预期工作?什么是“文件 162,140,数据输出:1.9G”?这些数字是描述输入还是输出?这对我们很重要吗?
  • 只是一个警告:if threads &lt; 0: return "Thead Count not specified" 本身并不能真正让人想要调试你的代码。负数表示“未指定”?在 Python 世界中,这就是语义 B$。默认情况下,不提供参数会给你一个很好的例外。并将错误消息作为字符串返回?那么,您将这个函数的返回值指定为None 还是错误消息?在尝试实现复杂的并发构造之前,您确实应该深入研究 Python 异常。

标签: python multithreading solr multiprocessing python-multiprocessing


【解决方案1】:

我在上面的 cmets 中给出了一些一般性的提示。我认为您的方法存在多个问题,在非常不同的抽象级别上。您也没有显示所有相关代码。

这个问题很可能是

  1. 在您用于从 solr 读取数据的方法中或在将读取数据提供给您的工作人员之前准备读取数据。
  2. 在您提出的用于在多个进程之间分配工作的架构中。
  3. 在您的日志基础设施中(正如您自己指出的那样)。
  4. 在您的分析方法中。

必须了解所有这些要点,鉴于问题的复杂性,这里肯定没有人能够为您确定确切的问题。

关于第 (3) 和 (4) 点:

如果您不确定日志文件的完整性,则应根据处理引擎的有效负载输出进行分析。我想说的是:日志文件可能只是您的数据处理的副产品。主要产品是您应该分析的东西。当然,正确处理日志也很重要。但这两个问题应该分开处理。

我对上面列表中第 (2) 点的贡献:

您基于multiprocessing 的解决方案特别令人怀疑的是您等待工人完成工作的方式。您似乎不确定应该通过哪种方法等待您的工人,因此您应用了三种不同的方法:

首先,您在 while 循环中监控队列的大小并等待它变为 0。这是一种非规范方法,实际上可能有效。

其次,你join()你的流程很奇怪:

for p in procs:
    logging.debug(str(os.getpid()) + " - Joining Process {}".format(p))
    p.join(1)

为什么你在这里定义一秒的超时,而不响应进程是否在该时间范围内实际终止?您应该真正加入一个进程,即等到它终止,或者您指定一个超时,如果该超时在进程完成之前到期,请特别处理这种情况。你的代码没有区分这些情况,所以p.join(1)就像写time.sleep(1)一样。

第三,你加入队列。

那么,在确保q.qsize() 返回0 并再等一秒之后,你真的认为加入队列很重要吗?它有什么不同吗? 这些方法中的一种应该就足够了,您需要考虑这些标准中的哪一个对您的问题最重要。也就是说,这些条件之一应该确定性地暗示其他两个条件。

所有这些看起来像是对多处理解决方案的快速而肮脏的破解,而您自己并不确定该解决方案的行为方式。我在处理并发架构时获得的最重要的见解之一:作为架构师,您必须 100% 了解通信和控制流在您的系统中是如何工作的。未正确监视和控制工作进程的状态很可能是您观察到的问题的根源。

【讨论】:

    【解决方案2】:

    我想通了,我听从 Jan-Philip 的建议,开始检查多进程/多线程进程的输出数据。结果发现,一个用 Solr 中的数据完成所有这些事情的对象在线程之间共享。我没有任何锁定机制,所以如果它有来自多个会话的混合数据,导致输出不一致。我通过为每个线程实例化一个新对象并且计数匹配来验证这一点。它有点慢,但仍然可行。

    谢谢

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-01-22
      • 1970-01-01
      • 2017-01-27
      • 2010-12-17
      • 1970-01-01
      • 2021-12-26
      • 2017-04-01
      • 1970-01-01
      相关资源
      最近更新 更多