最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

如何通過Python實現(xiàn)RabbitMQ延遲隊列

 更新時間:2020年11月28日 15:12:04   作者:Bge的博客  
這篇文章主要介紹了如何通過Python實現(xiàn)RabbitMQ延遲隊列,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下

最近在做一任務(wù)時,遇到需要延遲處理的數(shù)據(jù),最開始的做法是現(xiàn)將數(shù)據(jù)存儲在數(shù)據(jù)庫,然后寫個腳本,隔五分鐘掃描數(shù)據(jù)表再處理數(shù)據(jù),實際效果并不好。因為系統(tǒng)本身一直在用RabbitMQ做異步處理任務(wù)的中間件,所以想到是否可以利用RabbitMQ實現(xiàn)延遲隊列。功夫不負有心人,RabbitMQ雖然沒有現(xiàn)成可用的延遲隊列,但是可以利用其兩個重要特性來實現(xiàn)之:1、Time To Live(TTL)消息超時機制;2、Dead Letter Exchanges(DLX)死信隊列。下面將具體描述實現(xiàn)原理以及實現(xiàn)代

延遲隊列的基礎(chǔ)原理Time To Live(TTL)

RabbitMQ可以針對Queue設(shè)置x-expires 或者 針對Message設(shè)置 x-message-ttl,來控制消息的生存時間,如果超時(兩者同時設(shè)置以最先到期的時間為準),則消息變?yōu)閐ead letter(死信)
RabbitMQ消息的過期時間有兩種方法設(shè)置。

通過隊列(Queue)的屬性設(shè)置,隊列中所有的消息都有相同的過期時間。(本次延遲隊列采用的方案)對消息單獨設(shè)置,每條消息TTL可以不同。

如果同時使用,則消息的過期時間以兩者之間TTL較小的那個數(shù)值為準。消息在隊列的生存時間一旦超過設(shè)置的TTL值,就成為死信(dead letter)

Dead Letter Exchanges(DLX)

RabbitMQ的Queue可以配置x-dead-letter-exchange 和x-dead-letter-routing-key(可選)兩個參數(shù),如果隊列內(nèi)出現(xiàn)了dead letter,則按照這兩個參數(shù)重新路由轉(zhuǎn)發(fā)到指定的隊列。

  • x-dead-letter-exchange:出現(xiàn)死信(dead letter)之后將dead letter重新發(fā)送到指定exchange
  • x-dead-letter-routing-key:出現(xiàn)死信(dead letter)之后將dead letter重新按照指定的routing-key發(fā)送

隊列中出現(xiàn)死信(dead letter)的情況有:

  • 消息或者隊列的TTL過期。(延遲隊列利用的特性)
  • 隊列達到最大長度
  • 消息被消費端拒絕(basic.reject or basic.nack)并且requeue=false

綜合上面兩個特性,將隊列設(shè)置TTL規(guī)則,隊列TTL過期后消息會變成死信,然后利用DLX特性將其轉(zhuǎn)發(fā)到另外的交換機和隊列就可以被重新消費,達到延遲消費效果。

延遲隊列設(shè)計及實現(xiàn)(Python)

從上面描述,延遲隊列的實現(xiàn)大致分為兩步:

產(chǎn)生死信,有兩種方式Per-Message TTL和 Queue TTL,因為我的需求中是所有的消息延遲處理時間相同,所以本實現(xiàn)中采用 Queue TTL設(shè)置隊列的TTL,如果需要將隊列中的消息設(shè)置不同的延遲處理時間,則設(shè)置Per-Message TTL(官方文檔

設(shè)置死信的轉(zhuǎn)發(fā)規(guī)則,Dead Letter Exchanges設(shè)置方法(官方文檔

完整代碼如下:

"""
Created on Fri Aug 3 17:00:44 2018

@author: Bge
"""
import pika,json,logging
class RabbitMQClient:
  def __init__(self, conn_str='amqp://user:pwd@host:port/%2F'):
    self.exchange_type = "direct"
    self.connection_string = conn_str
    self.connection = pika.BlockingConnection(pika.URLParameters(self.connection_string))
    self.channel = self.connection.channel()
    self._declare_retry_queue() #RetryQueue and RetryExchange
    logging.debug("connection established")
  def close_connection(self):
    self.connection.close()
    logging.debug("connection closed")
  def declare_exchange(self, exchange):
    self.channel.exchange_declare(exchange=exchange,
                   exchange_type=self.exchange_type,
                   durable=True)
  def declare_queue(self, queue):
    self.channel.queue_declare(queue=queue,
                  durable=True,)
  def declare_delay_queue(self, queue,DLX='RetryExchange',TTL=60000):
    """
    創(chuàng)建延遲隊列
    :param TTL: ttl的單位是us,ttl=60000 表示 60s
    :param queue:
    :param DLX:死信轉(zhuǎn)發(fā)的exchange
    :return:
    """
    arguments={}
    if DLX:
      #設(shè)置死信轉(zhuǎn)發(fā)的exchange
      arguments[ 'x-dead-letter-exchange']=DLX
    if TTL:
      arguments['x-message-ttl']=TTL
    print(arguments)
    self.channel.queue_declare(queue=queue,
                  durable=True,
                  arguments=arguments)
  def _declare_retry_queue(self):
    """
    創(chuàng)建異常交換器和隊列,用于存放沒有正常處理的消息。
    :return:
    """
    self.channel.exchange_declare(exchange='RetryExchange',
                   exchange_type='fanout',
                   durable=True)
    self.channel.queue_declare(queue='RetryQueue',
                  durable=True)
    self.channel.queue_bind('RetryQueue', 'RetryExchange','RetryQueue')
  def publish_message(self,routing_key, msg,exchange='',delay=0,TTL=None):
    """
    發(fā)送消息到指定的交換器
    :param exchange: RabbitMQ交換器
    :param msg: 消息實體,是一個序列化的JSON字符串
    :return:
    """
    if delay==0:
      self.declare_queue(routing_key)
    else:
      self.declare_delay_queue(routing_key,TTL=TTL)
    if exchange!='':
      self.declare_exchange(exchange)
    self.channel.basic_publish(exchange=exchange,
                  routing_key=routing_key,
                  body=msg,
                  properties=pika.BasicProperties(
                    delivery_mode=2,
                    type=exchange
                  ))
    self.close_connection()
    print("message send out to %s" % exchange)
    logging.debug("message send out to %s" % exchange)
  def start_consume(self,callback,queue='#',delay=1):
    """
    啟動消費者,開始消費RabbitMQ中的消息
    :return:
    """
    if delay==1:
      queue='RetryQueue'
    else:
      self.declare_queue(queue)
    self.channel.basic_qos(prefetch_count=1)
    try:
      self.channel.basic_consume( # 消費消息
        callback, # 如果收到消息,就調(diào)用callback函數(shù)來處理消息
        queue=queue, # 你要從那個隊列里收消息
      )
      self.channel.start_consuming()
    except KeyboardInterrupt:
      self.stop_consuming()
  def stop_consuming(self):
    self.channel.stop_consuming()
    self.close_connection()
  def message_handle_successfully(channel, method):
    """
    如果消息處理正常完成,必須調(diào)用此方法,
    否則RabbitMQ會認為消息處理不成功,重新將消息放回待執(zhí)行隊列中
    :param channel: 回調(diào)函數(shù)的channel參數(shù)
    :param method: 回調(diào)函數(shù)的method參數(shù)
    :return:
    """
    channel.basic_ack(delivery_tag=method.delivery_tag)
  def message_handle_failed(channel, method):
    """
    如果消息處理失敗,應(yīng)該調(diào)用此方法,會自動將消息放入異常隊列
    :param channel: 回調(diào)函數(shù)的channel參數(shù)
    :param method: 回調(diào)函數(shù)的method參數(shù)
    :return:
    """
    channel.basic_reject(delivery_tag=method.delivery_tag, requeue=False)

發(fā)布消息代碼如下:

from MQ.RabbitMQ import RabbitMQClient
print("start program")
client = RabbitMQClient()
msg1 = '{"key":"value"}'
client.publish_message('test-delay',msg1,delay=1,TTL=10000)
print("message send out")

消費者代碼如下:

from MQ.RabbitMQ import RabbitMQClient
import json
print("start program")
client = RabbitMQClient()
def callback(ch, method, properties, body):
    msg = body.decode()
    print(msg)
    # 如果處理成功,則調(diào)用此消息回復ack,表示消息成功處理完成。
    RabbitMQClient.message_handle_successfully(ch, method)
queue_name = "RetryQueue"
client.start_consume(callback,queue_name,delay=0)

以上就是本文的全部內(nèi)容,希望對大家的學習有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

  • python虛擬環(huán)境創(chuàng)建的兩種方法

    python虛擬環(huán)境創(chuàng)建的兩種方法

    本文主要介紹了python虛擬環(huán)境創(chuàng)建的兩種方法,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2023-05-05
  • python實現(xiàn)360的字符顯示界面

    python實現(xiàn)360的字符顯示界面

    這篇文章主要介紹了python實現(xiàn)360的字符顯示界面示例,需要的朋友可以參考下
    2014-02-02
  • Python?列表和字典常踩坑即解決方案

    Python?列表和字典常踩坑即解決方案

    這篇文章介紹了?Python?列表和字典常踩坑即解決方案,主要針對python列表和字典得一些問題,提出了多種解決方案。需要的小伙伴可以參考一下
    2022-05-05
  • 使用beaker讓Facebook的Bottle框架支持session功能

    使用beaker讓Facebook的Bottle框架支持session功能

    這篇文章主要介紹了使用beaker讓Facebook的Bottle框架支持session功能,session在Python的Django等框架中內(nèi)置但在Bottle中并沒有被集成,需要的朋友可以參考下
    2015-04-04
  • 使用Python實現(xiàn)自動化辦公的代碼示例(郵件、excel)

    使用Python實現(xiàn)自動化辦公的代碼示例(郵件、excel)

    隨著技術(shù)的進步,Python 的高效性和易用性使其成為辦公自動化的強大工具,通過 Python,我們可以自動處理日常工作中的郵件、Excel 表格等任務(wù),從而大幅提升效率,本文將詳細介紹如何使用 Python 實現(xiàn)這些自動化功能,并附上關(guān)鍵代碼示例,需要的朋友可以參考下
    2025-01-01
  • django 讀取圖片到頁面實例

    django 讀取圖片到頁面實例

    這篇文章主要介紹了django 讀取圖片到頁面實例,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-03-03
  • python操作kafka實踐的示例代碼

    python操作kafka實踐的示例代碼

    這篇文章主要介紹了python操作kafka實踐的示例代碼,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2019-06-06
  • 基于TensorFlow中自定義梯度的2種方式

    基于TensorFlow中自定義梯度的2種方式

    今天小編就為大家分享一篇基于TensorFlow中自定義梯度的2種方式,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-02-02
  • 使用Python實現(xiàn)MapReduce的示例代碼

    使用Python實現(xiàn)MapReduce的示例代碼

    MapReduce是一個用于大規(guī)模數(shù)據(jù)處理的分布式計算模型,最初由Google工程師設(shè)計并實現(xiàn)的,Google已經(jīng)將完整的MapReduce論文公開發(fā)布了,本文給大家介紹了使用Python實現(xiàn)MapReduce的示例代碼,需要的朋友可以參考下
    2024-05-05
  • Python中sys模塊功能與用法實例詳解

    Python中sys模塊功能與用法實例詳解

    這篇文章主要介紹了Python中sys模塊功能與用法,結(jié)合實例形式詳細分析了Python sys模塊基本功能、原理、使用方法及操作注意事項,需要的朋友可以參考下
    2020-02-02

最新評論

如皋市| 大兴区| 浑源县| 偃师市| 乌鲁木齐市| 榆中县| 富裕县| 广水市| 石楼县| 大洼县| 阳谷县| 信阳市| 布尔津县| 铜鼓县| 九江市| 嘉义县| 威海市| 昌都县| 沁源县| 海淀区| 温州市| 青田县| 平定县| 伊吾县| 宁河县| 临江市| 米脂县| 西青区| 石渠县| 广德县| 西盟| 广安市| 泾阳县| 闽清县| 新竹市| 岗巴县| 龙江县| 邵阳县| 皋兰县| 康保县| 上栗县|