【问题标题】:Consume Redis messages with a pool of workers使用工人池使用 Redis 消息
【发布时间】:2016-01-28 14:22:24
【问题描述】:

我有一个 Redis 列表,其中发布者推送一些消息(JSON 序列化)。

另一方面,订阅者可以获取每个 JSON blob 并执行某些操作。最简单的方法是串行执行此操作。但我想让它快一点;我想维护一个工作进程池(多个消费者),每当有新消息到达时,检查池中是否有一个可以开始处理的“空闲”进程 我正在寻找以下基于池的版本

while not False:
    _, new_user = conn.blpop('queue:users')
    if not new_user:
        continue
    try:
        process_new_user(new_user, conn)
    except Exception as e:
        print e
    else:
        pass

但是我无法将其转换为使用 pythons multiprocessing.Pool 类的代码。文档也无济于事

【问题讨论】:

    标签: python redis message-queue publish-subscribe python-multiprocessing


    【解决方案1】:
    from multiprocessing import Pool
    
    
    pool = Pool()
    
    while 1:
        _, new_user = conn.blpop('queue:users')
    
        if not new_user:
            continue
    
        pool.apply_async(process_new_user, args=(new_user, conn))
    

    如果你想处理异常,你需要收集pool.apply_async返回的AsyncResult对象并检查它们的状态。

    如果您可以使用 Python 3,concurrent.futures 池允许在异步回调中处理结果,从而更容易检查作业退出状态。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-09-11
      • 2018-09-05
      • 2021-01-01
      • 2013-03-06
      • 2020-06-20
      • 1970-01-01
      • 2020-07-31
      • 1970-01-01
      相关资源
      最近更新 更多