文章详情

短信预约-IT技能 免费直播动态提醒

请输入下面的图形验证码

提交验证

短信预约提醒成功

python使用redis实现消息队列(异步)的实现完整例程

2023-01-18 11:00

关注

最近在用fastapi框架开发web后端,由于近几年python异步编程大火,fastapi凭借高性能也火了起来。本篇介绍了在异步环境下实现Redis消息队列的方法,代码可以直接拷贝到fastapi中使用。

安装相关库

pip install aioredis

消息队列实现及使用

我们使用redis的stream类型作为消息队列的载体

首先我们创建一个目录作为项目目录:works/

创建配置文件

在项目根目录下新建文件works/.env

在文件中写入

export APP_ENV=development

export REDIS_URL="192.168.70.130/"
export REDIS_USER=
export REDIS_PASSWORD=
export REDIS_HOST="192.168.70.130"
export REDIS_PORT=6379

代码实现

在项目目录下创建py文件works/main.py

import os
from dotenv import load_dotenv
import aioredis
import asyncio

load_dotenv()

class Redis():
  def __init__(self):
    """initialize connection """
    self.REDIS_URL = os.environ['REDIS_URL']
    self.REDIS_PASSWORD = os.environ['REDIS_PASSWORD']
    self.REDIS_USER = os.environ['REDIS_USER']
    self.connection_url = f"redis://{self.REDIS_USER}:{self.REDIS_PASSWORD}@{self.REDIS_URL}"
    self.REDIS_HOST = os.environ['REDIS_HOST']
    self.REDIS_PORT = os.environ['REDIS_PORT']
   
  async def create_connection(self):
    self.connection = aioredis.from_url(
      self.connection_url, db=0)

    return self.connection


class Producer:
  def __init__(self, redis_client):
    self.redis_client = redis_client

  async def add_to_stream(self, data: dict, stream_channel):
    """将一条数据添加到队列

    Args:
      data (dict): _description_
      stream_channel (_type_): _description_

    Returns:
      _type_: _description_
    """
    try:
      msg_id = await self.redis_client.xadd(name=stream_channel, id="*", fields=data)
      print(f"Message id {msg_id} added to {stream_channel} stream")
      return msg_id

    except Exception as e:
      raise Exception(f"Error sending msg to stream => {e}")

class StreamConsumer:
  def __init__(self, redis_client):
    self.redis_client = redis_client

  async def consume_stream(self, count: int, block: int, stream_channel):
    """读取队列中的消息,但是并不删除

    Args:
      count (int): _description_
      block (int): _description_
      stream_channel (_type_): _description_

    Returns:
      _type_: _description_
    """
    response = await self.redis_client.xread(
      streams={stream_channel: '0-0'}, count=count, block=block)

    return response

  async def delete_message(self, stream_channel, message_id):
    """成功消费数据后,调用此函数删除队列数据

    Args:
      stream_channel (_type_): _description_
      message_id (_type_): _description_
    """
    await self.redis_client.xdel(stream_channel, message_id)


async def main():
  redis_conn = await Redis().create_connection()
  produce = Producer(redis_conn)
  consumer = StreamConsumer(redis_conn)
  # 添加一个消息到队列中
  data = {'xiaoming4':123}
  await produce.add_to_stream(data=data,stream_channel='message_channel')
 
  # 从队列中拿出最新的1条数据
  data = await consumer.consume_stream(1,block=0,stream_channel='message_channel')
  print(data)
 
  # 轮询等待队列中的新消息
  response = await consumer.consume_stream(stream_channel="message_channel", count=1, block=0)
  if response:
    for stream, messagees in response:
      print('stream:',stream)
      for message in messagees:
        print('message: ',message)
        message_id = message[0]
        print('message_id: ',message_id)
        message_content = message[1]
        print('message_content: ',message_content)
        print('注意里面的键、值都变成了byte类型,需要进行解码:')
        message_content:dict
        print('message_content_decode: ',{k.decode('utf-8'):v.decode('utf-8') for k,v in message_content.items()})

  # 消费成功后删除队列中的消息
  await consumer.delete_message(
    stream_channel='message_channel',message_id=message_id
  )  

if __name__ == '__main__':
  asyncio.run(main())

非常简单好用,启动一下看看吧

到此这篇关于python使用redis实现消息队列(异步)的实现完整例程的文章就介绍到这了,更多相关python redis消息队列内容请搜索我们以前的文章或继续浏览下面的相关文章希望大家以后多多支持我们!

阅读原文内容投诉

免责声明:

① 本站未注明“稿件来源”的信息均来自网络整理。其文字、图片和音视频稿件的所属权归原作者所有。本站收集整理出于非商业性的教育和科研之目的,并不意味着本站赞同其观点或证实其内容的真实性。仅作为临时的测试数据,供内部测试之用。本站并未授权任何人以任何方式主动获取本站任何信息。

② 本站未注明“稿件来源”的临时测试数据将在测试完成后最终做删除处理。有问题或投稿请发送至: 邮箱/279061341@qq.com QQ/279061341

软考中级精品资料免费领

  • 历年真题答案解析
  • 备考技巧名师总结
  • 高频考点精准押题
  • 2024年上半年信息系统项目管理师第二批次真题及答案解析(完整版)

    难度     813人已做
    查看
  • 【考后总结】2024年5月26日信息系统项目管理师第2批次考情分析

    难度     354人已做
    查看
  • 【考后总结】2024年5月25日信息系统项目管理师第1批次考情分析

    难度     318人已做
    查看
  • 2024年上半年软考高项第一、二批次真题考点汇总(完整版)

    难度     435人已做
    查看
  • 2024年上半年系统架构设计师考试综合知识真题

    难度     224人已做
    查看

相关文章

发现更多好内容

猜你喜欢

AI推送时光机
位置:首页-资讯-后端开发
咦!没有更多了?去看看其它编程学习网 内容吧
首页课程
资料下载
问答资讯