在 Python 中使用带有 asyncio 的信号量

Jef*_*son 8 python semaphore python-asyncio

我试图限制使用信号量同时运行的异步函数的数量,但我无法让它工作。我的代码归结为:

import asyncio


async def send(i):

    print(f"starting {i}")
    await asyncio.sleep(4)
    print(f"ending {i}")


async def helper():
    async with asyncio.Semaphore(value=5):
        await asyncio.gather(*[
            send(1),
            send(2),
            send(3),
            send(4),
            send(5),
            send(6),
            send(7),
            send(8),
            send(9),
            send(10),
        ])


if __name__ == "__main__":
    loop = asyncio.get_event_loop()
    loop.run_until_complete(helper())
    loop.close()
Run Code Online (Sandbox Code Playgroud)

输出是:

starting 1
starting 2
starting 3
starting 4
starting 5
starting 6
starting 7
starting 8
starting 9
starting 10
ending 1
ending 2
ending 3
ending 4
ending 5
ending 6
ending 7
ending 8
ending 9
ending 10
Run Code Online (Sandbox Code Playgroud)

我希望并期望只有 5 个会同时运行,但是所有 10 个都同时启动和停止。我究竟做错了什么?

Art*_*rev 11

请找到下面的工作示例,请随时提出问题:

import asyncio


async def send(i: int, semaphore: asyncio.Semaphore):
    # to demonstrate that all tasks start nearly together
    print(f"Hello: {i}")
    # only two tasks can run code inside the block below simultaneously
    async with semaphore:
        print(f"starting {i}")
        await asyncio.sleep(4)
        print(f"ending {i}")


async def async_main():
    s = asyncio.Semaphore(value=2)
    await asyncio.gather(*[send(i, semaphore=s) for i in range(1, 11)])


if __name__ == "__main__":
    loop = asyncio.get_event_loop()
    loop.run_until_complete(async_main())
    loop.close()

Run Code Online (Sandbox Code Playgroud)

2023 年 8 月 18 日起的版本:

我看到很多人对如何使用感兴趣asyncio.Semaphore,我决定扩展我的答案。

新版本说明了如何将生产者-消费者模式与asyncio.Semaphore. 如果您想要非常简单的东西,您可以使用上面原始答案中的代码。如果您想要更强大的解决方案(允许限制asyncio.Tasks使用的数量),您可以使用这个更强大的解决方案。

import asyncio
from typing import List

CONSUMERS_NUMBER = 10  # workers/consumer number
TASKS_NUMBER = 20  # number of tasks to do


async def producer(tasks_to_do: List[int], q: asyncio.Queue) -> None:
    print(f"Producer started working!")
    for task in tasks_to_do:
        await q.put(task)  # put tasks to Queue

    # poison pill technique
    for _ in range(CONSUMERS_NUMBER):
        await q.put(None)  # put poison pill to all worker/consumers

    print("Producer finished working!")


async def consumer(
        consumer_name: str,
        q: asyncio.Queue,
        semaphore: asyncio.Semaphore,
) -> None:
    print(f"{consumer_name} started working!")
    while True:
        task = await q.get()

        if task is None:  # stop if poison pill was received
            break

        print(f"{consumer_name} took {task} from queue!")

        # number of tasks which could be processed simultaneously
        # is limited by semaphore
        async with semaphore:
            print(f"{consumer_name} started working with {task}!")
            await asyncio.sleep(4)
            print(f"{consumer_name} finished working with {task}!")

    print(f"{consumer_name} finished working!")


async def async_main() -> None:
    """Main entrypoint of async app."""
    tasks = [f"TheTask#{i + 1}" for i in range(TASKS_NUMBER)]
    q = asyncio.Queue(maxsize=2)
    s = asyncio.Semaphore(value=2)
    consumers = [
        consumer(
            consumer_name=f"Consumer#{i + 1}",
            q=q,
            semaphore=s,
        ) for i in range(CONSUMERS_NUMBER)
    ]
    await asyncio.gather(producer(tasks_to_do=tasks, q=q), *consumers)


if __name__ == "__main__":
    asyncio.run(async_main())


Run Code Online (Sandbox Code Playgroud)