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

相关文章

hdu1043(八数码问题,广搜 + hash(实现状态压缩) )

利用康拓展开将一个排列映射成一个自然数,然后就变成了普通的广搜题。 #include<iostream>#include<algorithm>#include<string>#include<stack>#include<queue>#include<map>#include<stdio.h>#include<stdlib.h>#include<ctype.h>#inclu

C++11第三弹:lambda表达式 | 新的类功能 | 模板的可变参数

🌈个人主页: 南桥几晴秋 🌈C++专栏: 南桥谈C++ 🌈C语言专栏: C语言学习系列 🌈Linux学习专栏: 南桥谈Linux 🌈数据结构学习专栏: 数据结构杂谈 🌈数据库学习专栏: 南桥谈MySQL 🌈Qt学习专栏: 南桥谈Qt 🌈菜鸡代码练习: 练习随想记录 🌈git学习: 南桥谈Git 🌈🌈🌈🌈🌈🌈🌈🌈🌈🌈🌈🌈🌈�

【C++】_list常用方法解析及模拟实现

相信自己的力量,只要对自己始终保持信心,尽自己最大努力去完成任何事,就算事情最终结果是失败了,努力了也不留遗憾。💓💓💓 目录   ✨说在前面 🍋知识点一:什么是list? •🌰1.list的定义 •🌰2.list的基本特性 •🌰3.常用接口介绍 🍋知识点二:list常用接口 •🌰1.默认成员函数 🔥构造函数(⭐) 🔥析构函数 •🌰2.list对象

【Prometheus】PromQL向量匹配实现不同标签的向量数据进行运算

✨✨ 欢迎大家来到景天科技苑✨✨ 🎈🎈 养成好习惯,先赞后看哦~🎈🎈 🏆 作者简介:景天科技苑 🏆《头衔》:大厂架构师,华为云开发者社区专家博主,阿里云开发者社区专家博主,CSDN全栈领域优质创作者,掘金优秀博主,51CTO博客专家等。 🏆《博客》:Python全栈,前后端开发,小程序开发,人工智能,js逆向,App逆向,网络系统安全,数据分析,Django,fastapi

让树莓派智能语音助手实现定时提醒功能

最初的时候是想直接在rasa 的chatbot上实现,因为rasa本身是带有remindschedule模块的。不过经过一番折腾后,忽然发现,chatbot上实现的定时,语音助手不一定会有响应。因为,我目前语音助手的代码设置了长时间无应答会结束对话,这样一来,chatbot定时提醒的触发就不会被语音助手获悉。那怎么让语音助手也具有定时提醒功能呢? 我最后选择的方法是用threading.Time

Android实现任意版本设置默认的锁屏壁纸和桌面壁纸(两张壁纸可不一致)

客户有些需求需要设置默认壁纸和锁屏壁纸  在默认情况下 这两个壁纸是相同的  如果需要默认的锁屏壁纸和桌面壁纸不一样 需要额外修改 Android13实现 替换默认桌面壁纸: 将图片文件替换frameworks/base/core/res/res/drawable-nodpi/default_wallpaper.*  (注意不能是bmp格式) 替换默认锁屏壁纸: 将图片资源放入vendo

C#实战|大乐透选号器[6]:实现实时显示已选择的红蓝球数量

哈喽,你好啊,我是雷工。 关于大乐透选号器在前面已经记录了5篇笔记,这是第6篇; 接下来实现实时显示当前选中红球数量,蓝球数量; 以下为练习笔记。 01 效果演示 当选择和取消选择红球或蓝球时,在对应的位置显示实时已选择的红球、蓝球的数量; 02 标签名称 分别设置Label标签名称为:lblRedCount、lblBlueCount

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略

Kubernetes PodSecurityPolicy:PSP能实现的5种主要安全策略 1. 特权模式限制2. 宿主机资源隔离3. 用户和组管理4. 权限提升控制5. SELinux配置 💖The Begin💖点点关注,收藏不迷路💖 Kubernetes的PodSecurityPolicy(PSP)是一个关键的安全特性,它在Pod创建之前实施安全策略,确保P

工厂ERP管理系统实现源码(JAVA)

工厂进销存管理系统是一个集采购管理、仓库管理、生产管理和销售管理于一体的综合解决方案。该系统旨在帮助企业优化流程、提高效率、降低成本,并实时掌握各环节的运营状况。 在采购管理方面,系统能够处理采购订单、供应商管理和采购入库等流程,确保采购过程的透明和高效。仓库管理方面,实现库存的精准管理,包括入库、出库、盘点等操作,确保库存数据的准确性和实时性。 生产管理模块则涵盖了生产计划制定、物料需求计划、

C++——stack、queue的实现及deque的介绍

目录 1.stack与queue的实现 1.1stack的实现  1.2 queue的实现 2.重温vector、list、stack、queue的介绍 2.1 STL标准库中stack和queue的底层结构  3.deque的简单介绍 3.1为什么选择deque作为stack和queue的底层默认容器  3.2 STL中对stack与queue的模拟实现 ①stack模拟实现