【问题标题】:Can I run cleanup code in daemon threads in python?我可以在 python 的守护线程中运行清理代码吗?
【发布时间】:2022-01-16 20:29:22
【问题描述】:

假设我有一些消费者守护线程,只要主线程将对象放到队列中,它们就会不断地从队列中获取对象,并对它们执行一些长时间的操作(几秒钟)。

问题在于,每当主线程完成时,守护线程在完成处理队列中剩余的任何内容之前都会被杀死。

我知道解决此问题的一种方法可能是等待守护线程完成处理队列中剩余的所有内容然后退出,但我很好奇是否有任何方法可以让守护线程“清理”在主线程退出时(即完成处理队列中剩余的任何内容)之后,没有明确地让主线程告诉守护线程开始清理。

这背后的动机是我制作了一个 python 包,它有一个日志处理程序类,当用户尝试记录某些内容时(例如使用logging.info("message")),它会将项目放入队列中,并且处理程序有一个守护线程发送网络上的日志。我希望守护线程可以自行清理,这样包的用户就不必手动确保他们的主线程等待日志处理程序完成处理。

最小的工作示例

# this code is in my package
class MyHandler(logging.Handler):
  def __init__(self, level):
    super().__init__(level=level)
    self.queue = Queue()
    self.thread = Thread(target=self.consume, daemon=True)
    self.thread.start()

  def emit(self, record):
    # This gets called whenever the user does logging.info, or similar
    self.queue.put(record)

  def consume(self):
    while True:
      record = self.queue.get()
      send(record) # send record over network, can take a few seconds (assume it never raises)
      self.queue.task_done()
# This is user's main code

# user will have to keep a reference to the handler for later. I want to avoid this.
my_handler = MyHandler()
# set up logging
logging.basicConfig(..., handlers=[..., my_handler])

# do some stuff...
logging.info("this will be sent over network")
# some more stuff...
logging.error("also sent over network")
# even more stuff

# before exiting must wait for handler to finish sending
# I don't want user to have to do this
my_hanler.queue.join()

【问题讨论】:

  • 守护线程在没有清理的情况下突然关闭,甚至没有finally 块或__exit__ 方法。它们不是工作的正确工具。

标签: python multithreading daemon


【解决方案1】:

您可以使用threading.main_thread.join(),它会像这样等到关机:

import threading
import logging
import queue

class MyHandler(logging.Handler):
  def __init__(self, level):
    super().__init__(level=level)
    self.queue = queue.Queue()
    self.thread = threading.Thread(target=self.consume)  # Not daemon

    # Shutdown thread
    threading.Thread(
        target=lambda: threading.main_thread().join() or self.queue.put(None)
        ).start()
        
    self.thread.start()

  def emit(self, record):
    # This gets called whenever the user does logging.info, or similar
    self.queue.put(record)

  def consume(self):
    while True:
      record = self.queue.get()
      if record is None:
          print("cleaning")
          return  # Cleanup
      print(record) # send record over network, can take a few seconds (assume it never raises)
      self.queue.task_done()

快速测试代码:

logging.getLogger().setLevel(logging.INFO)
logging.getLogger().addHandler(MyHandler(logging.INFO))
logging.info("Hello")
exit()

【讨论】:

  • 谢谢!非常优雅的解决方案!我很感激
【解决方案2】:

您可以使用atexit 等待守护线程关闭:

import queue, threading, time, logging, atexit

class MyHandler(logging.Handler):
  def __init__(self, level):
    super().__init__(level=level)
    self.queue = queue.Queue()
    self.thread = threading.Thread(target=self.consume, daemon=True)

    # Right before main thread exits, signal cleanup and wait until done
    atexit.register(lambda: self.queue.put(None) or self.thread.join())

    self.thread.start()

  def emit(self, record):
    # This gets called whenever the user does logging.info, or similar
    self.queue.put(record)

  def consume(self):
    while True:
      record = self.queue.get()
      if record is None:  # Cleanup requested
          print("cleaning")
          time.sleep(5)
          return
      print(record) # send record over network, can take a few seconds (assume it never raises)
      self.queue.task_done()

# Test code
logging.getLogger().setLevel(logging.INFO)
logging.getLogger().addHandler(MyHandler(logging.INFO))
logging.info("Hello")

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-09-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-11-04
    相关资源
    最近更新 更多