如何写消息队列消费脚本

wen 实用脚本 32

本文目录导读:

如何写消息队列消费脚本

  1. RabbitMQ 消费脚本
  2. Kafka 消费脚本
  3. Redis 队列消费脚本
  4. 通用消息队列消费框架
  5. 最佳实践建议

我来帮你编写消息队列消费脚本,以下给出几种常用消息队列的示例:

RabbitMQ 消费脚本

Python 版本 (使用 pika)

import pika
import json
import logging
from datetime import datetime
# 配置日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class RabbitMQConsumer:
    def __init__(self, host='localhost', port=5672, 
                 queue='task_queue', username='guest', password='guest'):
        self.host = host
        self.port = port
        self.queue = queue
        self.username = username
        self.password = password
    def connect(self):
        """建立连接"""
        credentials = pika.PlainCredentials(self.username, self.password)
        parameters = pika.ConnectionParameters(
            host=self.host,
            port=self.port,
            credentials=credentials
        )
        self.connection = pika.BlockingConnection(parameters)
        self.channel = self.connection.channel()
        self.channel.queue_declare(queue=self.queue, durable=True)
    def process_message(self, channel, method, properties, body):
        """处理消息的回调函数"""
        try:
            # 解析消息
            message = json.loads(body)
            logger.info(f"收到消息: {message}")
            # 业务处理
            # TODO: 实现具体的业务逻辑
            self.handle_message(message)
            # 手动确认消息
            channel.basic_ack(delivery_tag=method.delivery_tag)
            logger.info("消息处理完成")
        except Exception as e:
            logger.error(f"消息处理失败: {e}")
            # 拒绝消息,重新入队
            channel.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
    def handle_message(self, message):
        """业务处理逻辑(需根据实际需求修改)"""
        # 示例:处理订单消息
        if message.get('type') == 'order':
            order_id = message.get('order_id')
            logger.info(f"处理订单: {order_id}")
            # 执行业务操作
        elif message.get('type') == 'notification':
            # 处理通知
            pass
    def start_consuming(self):
        """启动消费"""
        self.connect()
        # 设置QoS,限制未确认消息数量
        self.channel.basic_qos(prefetch_count=1)
        # 注册回调函数
        self.channel.basic_consume(
            queue=self.queue,
            on_message_callback=self.process_message
        )
        logger.info(f"开始消费队列: {self.queue}")
        try:
            self.channel.start_consuming()
        except KeyboardInterrupt:
            logger.info("停止消费")
            self.stop()
    def stop(self):
        """停止消费"""
        if self.connection and not self.connection.is_closed:
            self.connection.close()
# 使用示例
if __name__ == "__main__":
    consumer = RabbitMQConsumer(
        host='localhost',
        queue='my_queue'
    )
    consumer.start_consuming()

Kafka 消费脚本

Python 版本 (使用 confluent-kafka)

from confluent_kafka import Consumer, KafkaError
import json
import logging
from datetime import datetime
logger = logging.getLogger(__name__)
class KafkaConsumer:
    def __init__(self, bootstrap_servers='localhost:9092',
                 group_id='my-group',
                 topics=['my-topic']):
        self.bootstrap_servers = bootstrap_servers
        self.group_id = group_id
        self.topics = topics
        # Kafka消费者配置
        self.config = {
            'bootstrap.servers': self.bootstrap_servers,
            'group.id': self.group_id,
            'auto.offset.reset': 'earliest',
            'enable.auto.commit': False,  # 手动提交
            'session.timeout.ms': 6000,
            'max.poll.interval.ms': 300000
        }
    def create_consumer(self):
        """创建消费者实例"""
        self.consumer = Consumer(self.config)
        self.consumer.subscribe(self.topics)
    def process_message(self, message):
        """处理单条消息"""
        try:
            data = json.loads(message.value().decode('utf-8'))
            logger.info(f"收到消息: {data}")
            # 业务处理
            # TODO: 实现具体业务逻辑
            self.handle_message(data)
            return True
        except Exception as e:
            logger.error(f"处理消息失败: {e}")
            return False
    def handle_message(self, message):
        """业务处理逻辑"""
        # 示例处理
        if isinstance(message, dict):
            logger.info(f"处理消息: {message.get('id')}")
    def start_consuming(self):
        """启动消费"""
        self.create_consumer()
        logger.info(f"开始消费主题: {self.topics}")
        try:
            while True:
                # 批量拉取消息
                msgs = self.consumer.consume(num_messages=100, timeout=1.0)
                if not msgs:
                    continue
                for msg in msgs:
                    if msg.error():
                        if msg.error().code() == KafkaError._PARTITION_EOF:
                            logger.warning("分区读取完毕")
                        else:
                            logger.error(f"消费错误: {msg.error()}")
                    else:
                        # 处理消息
                        success = self.process_message(msg)
                        # 手动提交偏移量
                        self.consumer.commit(msg)
        except KeyboardInterrupt:
            logger.info("停止消费")
        finally:
            self.consumer.close()
    def stop(self):
        """停止消费"""
        if self.consumer:
            self.consumer.close()
# 使用示例
if __name__ == "__main__":
    consumer = KafkaConsumer(
        bootstrap_servers='localhost:9092',
        topics=['my-topic']
    )
    consumer.start_consuming()

Redis 队列消费脚本

import redis
import json
import logging
from typing import Callable
import time
logger = logging.getLogger(__name__)
class RedisQueueConsumer:
    def __init__(self, host='localhost', port=6379, 
                 queue_name='my_queue', db=0):
        self.redis_client = redis.Redis(
            host=host, 
            port=port, 
            db=db,
            decode_responses=True
        )
        self.queue_name = queue_name
    def process_message(self, message: dict) -> bool:
        """处理消息(需根据实际需求修改)"""
        logger.info(f"处理消息: {message}")
        # TODO: 实现具体业务逻辑
        if message.get('type') == 'email':
            # 发送邮件
            pass
        elif message.get('type') == 'sms':
            # 发送短信
            pass
        return True
    def start_consuming(self, poll_interval=1.0, batch_size=10):
        """启动消费"""
        logger.info(f"开始消费Redis队列: {self.queue_name}")
        while True:
            try:
                # 批量获取消息
                pipeline = self.redis_client.pipeline()
                messages = []
                for _ in range(batch_size):
                    # 阻塞式获取消息,超时时间为poll_interval
                    message = self.redis_client.blpop(
                        self.queue_name, 
                        timeout=int(poll_interval)
                    )
                    if message:
                        queue_name, data = message
                        messages.append(json.loads(data))
                # 处理消息
                for message in messages:
                    try:
                        self.process_message(message)
                    except Exception as e:
                        logger.error(f"处理消息失败: {e}")
                        # 重新入队
                        self.redis_client.rpush(
                            self.queue_name, 
                            json.dumps(message)
                        )
                # 如果没有消息,等待
                if not messages:
                    time.sleep(poll_interval)
            except KeyboardInterrupt:
                logger.info("停止消费")
                break
            except Exception as e:
                logger.error(f"消费异常: {e}")
                time.sleep(5)  # 异常时等待后重试
    def stop(self):
        """停止消费"""
        self.redis_client.close()
# 使用示例
if __name__ == "__main__":
    consumer = RedisQueueConsumer(queue_name='task_queue')
    consumer.start_consuming()

通用消息队列消费框架

import abc
import logging
from typing import Any, Callable
logger = logging.getLogger(__name__)
class BaseMessageConsumer(abc.ABC):
    """消息队列消费基类"""
    def __init__(self, config: dict):
        self.config = config
        self.running = False
    @abc.abstractmethod
    def connect(self):
        """建立连接"""
        pass
    @abc.abstractmethod
    def consume(self, callback: Callable):
        """消费消息"""
        pass
    @abc.abstractmethod
    def acknowledge(self, message_id: str):
        """确认消息"""
        pass
    @abc.abstractmethod
    def reject(self, message_id: str):
        """拒绝消息"""
        pass
    def handle_message(self, message: Any) -> bool:
        """处理消息的默认方法"""
        logger.info(f"收到消息: {message}")
        return True
    def start(self, callback: Callable = None):
        """启动消费"""
        self.running = True
        self.connect()
        if callback:
            self.consume(callback)
        else:
            self.consume(self.handle_message)
    def stop(self):
        """停止消费"""
        self.running = False
        self.close()
    @abc.abstractmethod
    def close(self):
        """关闭连接"""
        pass
# 使用示例
class MyMessageConsumer(BaseMessageConsumer):
    def connect(self):
        # 实现连接逻辑
        pass
    def consume(self, callback):
        # 实现消费循环
        while self.running:
            # 获取消息
            message = self.get_message()
            if message:
                try:
                    success = callback(message)
                    if success:
                        self.acknowledge(message.id)
                    else:
                        self.reject(message.id)
                except Exception as e:
                    logger.error(f"处理失败: {e}")
                    self.reject(message.id)
    def get_message(self):
        # 实现获取消息逻辑
        pass
    def acknowledge(self, message_id):
        # 实现确认消息
        pass
    def reject(self, message_id):
        # 实现拒绝消息
        pass
    def close(self):
        # 实现关闭连接
        pass

最佳实践建议

消息处理策略

# 重试机制
def retry_on_failure(max_retries=3, delay=1.0):
    def decorator(func):
        def wrapper(message, *args, **kwargs):
            for attempt in range(max_retries):
                try:
                    return func(message, *args, **kwargs)
                except Exception as e:
                    if attempt < max_retries - 1:
                        time.sleep(delay * (attempt + 1))
                        continue
                    else:
                        raise e
            return None
        return wrapper
    return decorator
@retry_on_failure(max_retries=3)
def handle_order_message(message):
    # 处理订单消息
    pass

错误处理与监控

import sentry_sdk
class MonitoredConsumer:
    def process_message_with_monitoring(self, message):
        try:
            # 记录处理开始
            logger.info(f"开始处理消息: {message.id}")
            # 处理消息
            self.process_message(message)
            # 记录成功
            metrics.counter('message.processed.success').inc()
        except Exception as e:
            # 记录失败
            logger.error(f"处理失败: {e}")
            metrics.counter('message.processed.failure').inc()
            # 发送告警
            sentry_sdk.capture_exception(e)
            raise

性能优化

# 批量处理
def batch_consume_messages(consumer, batch_size=100):
    messages = []
    for _ in range(batch_size):
        message = consumer.get_message()
        if message:
            messages.append(message)
        else:
            break
    if messages:
        # 批量处理
        for message in messages:
            process_message(message)
        # 批量确认
        consumer.batch_acknowledge(messages)

选择哪种消息队列取决于你的具体需求,包括性能要求、可靠性要求、以及系统的其他组件等技术栈匹配情况。

抱歉,评论功能暂时关闭!