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

相关文章

Python使用python-can实现合并BLF文件

《Python使用python-can实现合并BLF文件》python-can库是Python生态中专注于CAN总线通信与数据处理的强大工具,本文将使用python-can为BLF文件合并提供高效灵活... 目录一、python-can 库:CAN 数据处理的利器二、BLF 文件合并核心代码解析1. 基础合

Python使用OpenCV实现获取视频时长的小工具

《Python使用OpenCV实现获取视频时长的小工具》在处理视频数据时,获取视频的时长是一项常见且基础的需求,本文将详细介绍如何使用Python和OpenCV获取视频时长,并对每一行代码进行深入解析... 目录一、代码实现二、代码解析1. 导入 OpenCV 库2. 定义获取视频时长的函数3. 打开视频文

golang版本升级如何实现

《golang版本升级如何实现》:本文主要介绍golang版本升级如何实现问题,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录golanwww.chinasem.cng版本升级linux上golang版本升级删除golang旧版本安装golang最新版本总结gola

SpringBoot中SM2公钥加密、私钥解密的实现示例详解

《SpringBoot中SM2公钥加密、私钥解密的实现示例详解》本文介绍了如何在SpringBoot项目中实现SM2公钥加密和私钥解密的功能,通过使用Hutool库和BouncyCastle依赖,简化... 目录一、前言1、加密信息(示例)2、加密结果(示例)二、实现代码1、yml文件配置2、创建SM2工具

Mysql实现范围分区表(新增、删除、重组、查看)

《Mysql实现范围分区表(新增、删除、重组、查看)》MySQL分区表的四种类型(范围、哈希、列表、键值),主要介绍了范围分区的创建、查询、添加、删除及重组织操作,具有一定的参考价值,感兴趣的可以了解... 目录一、mysql分区表分类二、范围分区(Range Partitioning1、新建分区表:2、分

MySQL 定时新增分区的实现示例

《MySQL定时新增分区的实现示例》本文主要介绍了通过存储过程和定时任务实现MySQL分区的自动创建,解决大数据量下手动维护的繁琐问题,具有一定的参考价值,感兴趣的可以了解一下... mysql创建好分区之后,有时候会需要自动创建分区。比如,一些表数据量非常大,有些数据是热点数据,按照日期分区MululbU

MySQL中查找重复值的实现

《MySQL中查找重复值的实现》查找重复值是一项常见需求,比如在数据清理、数据分析、数据质量检查等场景下,我们常常需要找出表中某列或多列的重复值,具有一定的参考价值,感兴趣的可以了解一下... 目录技术背景实现步骤方法一:使用GROUP BY和HAVING子句方法二:仅返回重复值方法三:返回完整记录方法四:

IDEA中新建/切换Git分支的实现步骤

《IDEA中新建/切换Git分支的实现步骤》本文主要介绍了IDEA中新建/切换Git分支的实现步骤,通过菜单创建新分支并选择是否切换,创建后在Git详情或右键Checkout中切换分支,感兴趣的可以了... 前提:项目已被Git托管1、点击上方栏Git->NewBrancjsh...2、输入新的分支的

Python实现对阿里云OSS对象存储的操作详解

《Python实现对阿里云OSS对象存储的操作详解》这篇文章主要为大家详细介绍了Python实现对阿里云OSS对象存储的操作相关知识,包括连接,上传,下载,列举等功能,感兴趣的小伙伴可以了解下... 目录一、直接使用代码二、详细使用1. 环境准备2. 初始化配置3. bucket配置创建4. 文件上传到os

关于集合与数组转换实现方法

《关于集合与数组转换实现方法》:本文主要介绍关于集合与数组转换实现方法,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录1、Arrays.asList()1.1、方法作用1.2、内部实现1.3、修改元素的影响1.4、注意事项2、list.toArray()2.1、方