File tree Expand file tree Collapse file tree 1 file changed +2
-1
lines changed
taskiq_aio_pika/taskiq/brokers Expand file tree Collapse file tree 1 file changed +2
-1
lines changed Original file line number Diff line number Diff line change @@ -14,7 +14,7 @@ class AioPikaBroker(AsyncBroker):
14
14
def __init__ (
15
15
self ,
16
16
result_backend : Optional [AsyncResultBackend [_T ]] = None ,
17
- qos : int = 1 ,
17
+ qos : int = 10 ,
18
18
loop : Optional [AbstractEventLoop ] = None ,
19
19
max_channel_pool_size : int = 2 ,
20
20
max_connection_pool_size : int = 10 ,
@@ -77,6 +77,7 @@ async def kick(self, message: TaskiqMessage) -> None:
77
77
78
78
async def listen (self ) -> AsyncGenerator [TaskiqMessage , None ]:
79
79
async with self .channel_pool .acquire () as channel :
80
+ await channel .set_qos (prefetch_count = self .qos )
80
81
queue = await channel .get_queue (self .queue_name , ensure = False )
81
82
async with queue .iterator () as queue_iter :
82
83
async for rmq_message in queue_iter :
You can’t perform that action at this time.
0 commit comments