代码之家  ›  专栏  ›  技术社区  ›  ddd

如何优雅地杀死侦听消息队列的线程

  •  1
  • ddd  · 技术社区  · 7 年前

    在我的Python应用程序中,我有一个函数,它使用来自amazonsqsfio队列的消息。

    def consume_msgs():
        sqs = boto3.client('sqs',
                       region_name='us-east-1',
                       aws_access_key_id=AWS_ACCESS_KEY_ID,
                       aws_secret_access_key=AWS_SECRET_ACCESS_KEY)
        print('STARTING WORKER listening on {}'.format(QUEUE_URL))
        while 1:
            response = sqs.receive_message(
                QueueUrl=QUEUE_URL,
                MaxNumberOfMessages=1,
                WaitTimeSeconds=10,
            )
            messages = response.get('Messages', [])
            for message in messages:
                try:
                    print('{} > {}'.format(threading.currentThread().getName(), message.get('Body')))
                    body = json.loads(message.get('Body'))
                    sqs.delete_message(QueueUrl=QUEUE_URL, ReceiptHandle=message.get('ReceiptHandle'))
    
                except Exception as e:
                    print('Exception in worker > ', e)
                    sqs.delete_message(QueueUrl=QUEUE_URL, ReceiptHandle=message.get('ReceiptHandle'))
    
        time.sleep(10)
    

    为了扩大规模,我使用多线程处理消息。

    if __name__ == '__main__:
    
        for i in range(3):
            t = threading.Thread(target=consume_msgs, name='worker-%s' % i)
            t.setDaemon(True)
            t.start()
        while True:
            print('Waiting')
            time.sleep(5)
    

    应用程序作为服务运行。如果需要部署新版本,则必须重新启动。当主进程被终止时,线程有优美的方式存在吗?它们不是突然终止线程,而是先完成当前消息,然后停止接收下一个消息。

    1 回复  |  直到 7 年前
        1
  •  2
  •   Ondrej K.    7 年前

    既然你的线程一直在循环,你就不能 join 他们,但你需要向他们发出信号,它的时间打破循环,以便能够做到这一点。这个 docs 提示可能有用:

    守护进程线程在关闭时突然停止。它们的资源(如打开的文件、数据库事务等)可能无法正确释放。如果希望线程正常停止,请使其非守护进程,并使用适当的信令机制,如 Event .

    有了这个,我把下面的例子放在一起,希望能有所帮助:

    from threading import Thread, Event
    from time import sleep
    
    def fce(ident, wrap_up_event):
        cnt = 0
        while True:
            print(f"{ident}: {cnt}", wrap_up_event.is_set())
            sleep(3)
            cnt += 1
            if wrap_up_event.is_set():
                break
        print(f"{ident}: Wrapped up")
    
    if __name__ == '__main__':
        wanna_exit = Event()
        for i in range(3):
            t = Thread(target=fce, args=(i, wanna_exit))
            t.start()
        sleep(5)
        wanna_exit.set()
    

    将单个事件实例传递给 fce 如果事件设置为 True . 在退出脚本之前,我们将此事件设置为 真的 从控制线程。由于线程不再标记为守护进程线程,因此我们不必显式地 参加 他们。

    根据您到底想如何关闭脚本,您需要处理传入的信号( SIGTERM 或许)或者 KeyboardInterrupt 例外情况 SIGINT . 在退出之前执行清理工作,其机制保持不变。除了不让python立即停止执行之外,还需要让线程知道它们不应该重新进入循环并等待它们加入。


    这个 西格特 有点简单,因为它公开为一个python异常,您可以对“main”位执行以下操作:

    if __name__ == '__main__':
        wanna_exit = Event()
        for i in range(3):
            t = Thread(target=fce, args=(i, wanna_exit))
            t.start()
        try:
            while True:
                sleep(5)
                print('Waiting')
        except KeyboardInterrupt:
            pass
        wanna_exit.set()
    

    你当然可以送 西格特 到一个进程 kill 不仅仅是控制终端。