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

相关文章

Spring Boot 实现 IP 限流的原理、实践与利弊解析

《SpringBoot实现IP限流的原理、实践与利弊解析》在SpringBoot中实现IP限流是一种简单而有效的方式来保障系统的稳定性和可用性,本文给大家介绍SpringBoot实现IP限... 目录一、引言二、IP 限流原理2.1 令牌桶算法2.2 漏桶算法三、使用场景3.1 防止恶意攻击3.2 控制资源

springboot下载接口限速功能实现

《springboot下载接口限速功能实现》通过Redis统计并发数动态调整每个用户带宽,核心逻辑为每秒读取并发送限定数据量,防止单用户占用过多资源,确保整体下载均衡且高效,本文给大家介绍spring... 目录 一、整体目标 二、涉及的主要类/方法✅ 三、核心流程图解(简化) 四、关键代码详解1️⃣ 设置

Nginx 配置跨域的实现及常见问题解决

《Nginx配置跨域的实现及常见问题解决》本文主要介绍了Nginx配置跨域的实现及常见问题解决,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来... 目录1. 跨域1.1 同源策略1.2 跨域资源共享(CORS)2. Nginx 配置跨域的场景2.1

Python中提取文件名扩展名的多种方法实现

《Python中提取文件名扩展名的多种方法实现》在Python编程中,经常会遇到需要从文件名中提取扩展名的场景,Python提供了多种方法来实现这一功能,不同方法适用于不同的场景和需求,包括os.pa... 目录技术背景实现步骤方法一:使用os.path.splitext方法二:使用pathlib模块方法三

CSS实现元素撑满剩余空间的五种方法

《CSS实现元素撑满剩余空间的五种方法》在日常开发中,我们经常需要让某个元素占据容器的剩余空间,本文将介绍5种不同的方法来实现这个需求,并分析各种方法的优缺点,感兴趣的朋友一起看看吧... css实现元素撑满剩余空间的5种方法 在日常开发中,我们经常需要让某个元素占据容器的剩余空间。这是一个常见的布局需求

HTML5 getUserMedia API网页录音实现指南示例小结

《HTML5getUserMediaAPI网页录音实现指南示例小结》本教程将指导你如何利用这一API,结合WebAudioAPI,实现网页录音功能,从获取音频流到处理和保存录音,整个过程将逐步... 目录1. html5 getUserMedia API简介1.1 API概念与历史1.2 功能与优势1.3

Java实现删除文件中的指定内容

《Java实现删除文件中的指定内容》在日常开发中,经常需要对文本文件进行批量处理,其中,删除文件中指定内容是最常见的需求之一,下面我们就来看看如何使用java实现删除文件中的指定内容吧... 目录1. 项目背景详细介绍2. 项目需求详细介绍2.1 功能需求2.2 非功能需求3. 相关技术详细介绍3.1 Ja

使用Python和OpenCV库实现实时颜色识别系统

《使用Python和OpenCV库实现实时颜色识别系统》:本文主要介绍使用Python和OpenCV库实现的实时颜色识别系统,这个系统能够通过摄像头捕捉视频流,并在视频中指定区域内识别主要颜色(红... 目录一、引言二、系统概述三、代码解析1. 导入库2. 颜色识别函数3. 主程序循环四、HSV色彩空间详解

PostgreSQL中MVCC 机制的实现

《PostgreSQL中MVCC机制的实现》本文主要介绍了PostgreSQL中MVCC机制的实现,通过多版本数据存储、快照隔离和事务ID管理实现高并发读写,具有一定的参考价值,感兴趣的可以了解一下... 目录一 MVCC 基本原理python1.1 MVCC 核心概念1.2 与传统锁机制对比二 Postg

SpringBoot整合Flowable实现工作流的详细流程

《SpringBoot整合Flowable实现工作流的详细流程》Flowable是一个使用Java编写的轻量级业务流程引擎,Flowable流程引擎可用于部署BPMN2.0流程定义,创建这些流程定义的... 目录1、流程引擎介绍2、创建项目3、画流程图4、开发接口4.1 Java 类梳理4.2 查看流程图4