Как с этим справиться, если количество сообщений меньше количества потоков?Python

Программы на Python
Anonymous
Как с этим справиться, если количество сообщений меньше количества потоков?

Сообщение Anonymous »

Извините за мой плохой английский.
При использовании многопоточности для обработки сообщений из очереди сообщений RabbitMQ, как это следует обрабатывать, если количество сообщений меньше количества потоки?
В моем коде для обработки можно использовать максимум 5 потоков. Если в очереди сообщений 13 URL-адресов, то по завершении работы программы 3 URL-адреса останутся необработанными. Если в очереди сообщений 12 URL-адресов, то 2 URL-адреса останутся необработанными.
Я плохо разбираюсь в RabbitMQ, но думаю, проблема связана с логикой Код RabbitMQ.

Код: Выделить всё

async def processTask(message: aio_pika.abc.AbstractIncomingMessage):
try:
messagesList.append(message)
# when message number =5
if len(messagesList) == 5:
messages_to_process = messagesList[:]

max_workers = 5 if len(messages_to_process) >= 5 else len(messages_to_process)
with ThreadPoolExecutor(max_workers=max_workers) as pool:
futures = {pool.submit(pageCollect, messageItem): messageItem for messageItem in
messages_to_process}
for future in as_completed(futures):
messageItem = futures[future]
try:
resultDict = future.result()
if resultDict:
await messageItem.ack()

logger.info(f'{messageItem.delivery_tag}: {resultDict}')
else:
await messageItem.reject(requeue=True)
except Exception as e:
logger.error(f'{messageItem.delivery_tag}: {e}')
await messageItem.reject(requeue=True)
linkList.clear()
messagesList.clear()

# when message number < 5
remaining_messages = queue.declaration_result.message_count
if remaining_messages > 0:
messages_to_process = messagesList[:]
max_workers = len(messages_to_process)
with ThreadPoolExecutor(max_workers=max_workers) as pool:
futures = {pool.submit(pageCollect, messageItem): messageItem for messageItem in messagesList}
for future in as_completed(futures):
messageItem = messages_to_process.pop(0)
try:
resultDict = future.result()
if resultDict:
await messageItem.ack()
logger.info(f'{messageItem.delivery_tag}: {resultDict}')
else:
await messageItem.reject(requeue=True)
except Exception as e:
logger.error(f'{messageItem.delivery_tag}: {e}')
await messageItem.reject(requeue=True)
messagesList.clear()

await asyncio.sleep(0.1)

except Exception as e:
print(e)

async def main(loop):
try:
# connect
connection = await aio_pika.connect_robust(host='XX.XX.X.XX', port=5672, login='admin', password='admin',
virtualhost='my_vhost', loop=loop)

global channel
channel = await connection.channel()
# Will take no more than 10 messages in advance
await channel.set_qos() #prefetch_count=5
crawler_exchange = await channel.declare_exchange(name='crawler_exchange', type='fanout')

queueName = "myqueue"
global queue
queue = await channel.declare_queue(queueName, durable=True)
await queue.bind(crawler_exchange, routing_key="myqueue")

rstqueueName = "Result"
global resultQuequ
resultQuequ = await channel.declare_queue(rstqueueName, durable=True)

global resultExchange
resultExchange = await channel.declare_exchange(name='resultExchange', type='direct')

await resultQuequ.bind(resultExchange, routing_key="allResult")

# get message
await queue.consume(processTask)
logger.info(f"Waiting for messages at {queue.name}.  To exit press CTRL+C")
return connection

except Exception as e:
logger.error(f"failed: {e}")
logger.error(traceback.format_exc())

if __name__ == '__main__':
loop = asyncio.get_event_loop()
connection = loop.run_until_complete(main(loop))

try:
loop.run_forever()
except KeyboardInterrupt:
logger.info("Received exit signal")
finally:
loop.run_until_complete(connection.close())
loop.close()

Я пытался изменить свой код с помощью gpt4 и claude3, но их код работает неправильно. Как мне решить проблему? Спасибо.
Что мне нужно:
Когда количество сообщений в очереди сообщений >= 5, должно быть 5 рабочих потоков. Когда количество сообщений < 5, количество рабочих потоков должно равняться количеству сообщений. Это может повысить скорость обработки и гарантировать, что в очереди сообщений не останется необработанных URL-адресов.

Подробнее здесь: https://stackoverflow.com/questions/782 ... mber-of-th

Вернуться в «Python»