aio_pika篇---实现收发功能

2023-11-09 15:50
文章标签 实现 功能 收发 pika aio

本文主要是介绍aio_pika篇---实现收发功能,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

aip_pika篇—实现收发功能

发送

publisher.py

import asynciofrom aio_pika import DeliveryMode, ExchangeType, Message, connectasync def main() -> None:# Perform connectionconnection = await connect(host='127.0.0.1',port=5672,login='ai.litter',password='r7n2ApE2yYk3yVz',virtualhost='sli')async with connection:# Creating a channelchannel = await connection.channel()logs_exchange = await channel.declare_exchange("ai.litter", ExchangeType.DIRECT,durable=True)# Sending the messagefor i in range(1000):message_body = 'hello  - {}'.format(str(i)).encode()message = Message(body=message_body,delivery_mode=DeliveryMode.PERSISTENT,)await asyncio.sleep(1)# routing_key = "hello"routing_key = "pool_queue"await logs_exchange.publish(message, routing_key=routing_key)print(f" [x] Sent {message.body!r}")if __name__ == "__main__":asyncio.run(main())

在这里插入图片描述

接收并发送

import aio_pika
from aio_pika import ExchangeType,Message
from aio_pika.abc import AbstractRobustConnection, AbstractIncomingMessage
from aio_pika.pool import Pool
from utils import receive_taskasync def main() -> None:loop = asyncio.get_event_loop()async def get_connection() -> AbstractRobustConnection:# return await aio_pika.connect_robust("amqp://guest:guest@localhost/")return await aio_pika.connect_robust(host='127.0.0.1', port=5672, login='ai.litter',password='r7n2ApE2yYk3yVz', virtualhost='slife')connection_pool: Pool = Pool(get_connection, max_size=2, loop=loop)async def get_channel() -> aio_pika.Channel:async with connection_pool.acquire() as connection:return await connection.channel()channel_pool: Pool = Pool(get_channel, max_size=10, loop=loop)async def consume() -> None:async with channel_pool.acquire() as channel:await channel.set_qos(20)direct_exchange = await channel.declare_exchange("ai.litter", ExchangeType.DIRECT, durable=True)queue_name = "pool_queue"queue = await channel.declare_queue(queue_name, durable=False, auto_delete=False,)await queue.bind(direct_exchange, routing_key=queue_name)async with queue.iterator() as queue_iter:message: AbstractIncomingMessageasync for message in queue_iter:try:print('task received, handling')print(str(message.body))await receive_task(message.body, publish_func=publish)except Exception as e:print('message nacked, exception=', e)await message.nack(requeue=False)else:print('task finished')try:await message.ack()except:await channel.reopen()async def publish(message: bytes, queue_name: str) -> None:async with channel_pool.acquire() as channel:# queue_name = "test_queue"routing_key = "test_queue"# Declaring exchangeexchange = await channel.declare_exchange("direct", auto_delete=True)# Declaring queuequeue = await channel.declare_queue(queue_name, auto_delete=True)# Binding queueawait queue.bind(exchange, routing_key)await exchange.publish(Message(# bytes("wwwwwwwwwwwwww", "utf-8"),message,# content_type="text/plain",# headers={"foo": "bar"},),routing_key,)async with connection_pool, channel_pool:task = loop.create_task(consume())print('amqp consumer created, waiting for task...')# await asyncio.wait([publish('hello world -- {}'.format(str(i)).encode(), queue_name) for i in range(5)])await taskif __name__ == "__main__":asyncio.run(main())import asyncio

在这里插入图片描述

查看收到的内容

utils.py

from typing import Callable, Coroutine, Any
import asyncio
from asyncio import queuesasync def receive_task(body: bytes, publish_func: Callable[[bytes, str], Coroutine[Any, Any, None]]):print("aaa ----->  " + str(body))await asyncio.wait([publish_func('good Job -- {}'.format(str(body)).encode(), 'hello')])

这篇关于aio_pika篇---实现收发功能的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/377097

相关文章

Oracle查询优化之高效实现仅查询前10条记录的方法与实践

《Oracle查询优化之高效实现仅查询前10条记录的方法与实践》:本文主要介绍Oracle查询优化之高效实现仅查询前10条记录的相关资料,包括使用ROWNUM、ROW_NUMBER()函数、FET... 目录1. 使用 ROWNUM 查询2. 使用 ROW_NUMBER() 函数3. 使用 FETCH FI

Python脚本实现自动删除C盘临时文件夹

《Python脚本实现自动删除C盘临时文件夹》在日常使用电脑的过程中,临时文件夹往往会积累大量的无用数据,占用宝贵的磁盘空间,下面我们就来看看Python如何通过脚本实现自动删除C盘临时文件夹吧... 目录一、准备工作二、python脚本编写三、脚本解析四、运行脚本五、案例演示六、注意事项七、总结在日常使用

Java实现Excel与HTML互转

《Java实现Excel与HTML互转》Excel是一种电子表格格式,而HTM则是一种用于创建网页的标记语言,虽然两者在用途上存在差异,但有时我们需要将数据从一种格式转换为另一种格式,下面我们就来看看... Excel是一种电子表格格式,广泛用于数据处理和分析,而HTM则是一种用于创建网页的标记语言。虽然两

Java中Springboot集成Kafka实现消息发送和接收功能

《Java中Springboot集成Kafka实现消息发送和接收功能》Kafka是一个高吞吐量的分布式发布-订阅消息系统,主要用于处理大规模数据流,它由生产者、消费者、主题、分区和代理等组件构成,Ka... 目录一、Kafka 简介二、Kafka 功能三、POM依赖四、配置文件五、生产者六、消费者一、Kaf

使用Python实现在Word中添加或删除超链接

《使用Python实现在Word中添加或删除超链接》在Word文档中,超链接是一种将文本或图像连接到其他文档、网页或同一文档中不同部分的功能,本文将为大家介绍一下Python如何实现在Word中添加或... 在Word文档中,超链接是一种将文本或图像连接到其他文档、网页或同一文档中不同部分的功能。通过添加超

windos server2022里的DFS配置的实现

《windosserver2022里的DFS配置的实现》DFS是WindowsServer操作系统提供的一种功能,用于在多台服务器上集中管理共享文件夹和文件的分布式存储解决方案,本文就来介绍一下wi... 目录什么是DFS?优势:应用场景:DFS配置步骤什么是DFS?DFS指的是分布式文件系统(Distr

NFS实现多服务器文件的共享的方法步骤

《NFS实现多服务器文件的共享的方法步骤》NFS允许网络中的计算机之间共享资源,客户端可以透明地读写远端NFS服务器上的文件,本文就来介绍一下NFS实现多服务器文件的共享的方法步骤,感兴趣的可以了解一... 目录一、简介二、部署1、准备1、服务端和客户端:安装nfs-utils2、服务端:创建共享目录3、服

C#使用yield关键字实现提升迭代性能与效率

《C#使用yield关键字实现提升迭代性能与效率》yield关键字在C#中简化了数据迭代的方式,实现了按需生成数据,自动维护迭代状态,本文主要来聊聊如何使用yield关键字实现提升迭代性能与效率,感兴... 目录前言传统迭代和yield迭代方式对比yield延迟加载按需获取数据yield break显式示迭

Python实现高效地读写大型文件

《Python实现高效地读写大型文件》Python如何读写的是大型文件,有没有什么方法来提高效率呢,这篇文章就来和大家聊聊如何在Python中高效地读写大型文件,需要的可以了解下... 目录一、逐行读取大型文件二、分块读取大型文件三、使用 mmap 模块进行内存映射文件操作(适用于大文件)四、使用 pand

python实现pdf转word和excel的示例代码

《python实现pdf转word和excel的示例代码》本文主要介绍了python实现pdf转word和excel的示例代码,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价... 目录一、引言二、python编程1,PDF转Word2,PDF转Excel三、前端页面效果展示总结一