本文目录导读:

我来帮你编写消息队列消费脚本,以下给出几种常用消息队列的示例:
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)
选择哪种消息队列取决于你的具体需求,包括性能要求、可靠性要求、以及系统的其他组件等技术栈匹配情况。