本文目录导读:

在PHP项目中监控和清理死信消息(Dead Letter Messages),通常涉及消息队列系统(如RabbitMQ、Kafka、Redis等),以下是系统的监控和清理方案:
常见死信场景
// 死信常见原因 - 消息被消费者拒绝(Reject)且 requeue=false - 消息TTL过期 - 队列达到最大长度 - 消费者处理失败次数超限
监控方案
1 使用消息队列管理工具
RabbitMQ 管理界面
# 查看死信队列 rabbitmqctl list_queues name messages messages_unacknowledged messages_ready # 查看死信交换器 rabbitmqctl list_exchanges name type
2 自定义监控脚本
<?php
class DeadLetterMonitor
{
private $connection;
public function __construct()
{
$this->connection = new AMQPConnection([
'host' => 'localhost',
'port' => 5672,
'vhost' => '/',
'login' => 'guest',
'password' => 'guest'
]);
}
// 监控死信队列大小
public function monitorDeadLetterQueue($deadQueueName)
{
$channel = new AMQPChannel($this->connection);
$queue = new AMQPQueue($channel);
$queue->setName($deadQueueName);
return [
'name' => $deadQueueName,
'message_count' => $queue->declareQueue(),
'consumer_count' => $queue->getConsumerCount()
];
}
// 发送告警
public function checkAndAlert($threshold = 100)
{
$metrics = $this->monitorDeadLetterQueue('dead_letter_queue');
if ($metrics['message_count'] > $threshold) {
// 发送告警通知
$this->sendAlert("死信队列消息数: {$metrics['message_count']}");
// 记录日志
$this->logDeadLetter($metrics);
}
}
private function sendAlert($message)
{
// 实现邮件、短信、Slack等通知
// mail(), Http client etc.
}
private function logDeadLetter($metrics)
{
$log = date('Y-m-d H:i:s') . " | " .
json_encode($metrics) . PHP_EOL;
file_put_contents('/var/log/dead_letter_monitor.log', $log, FILE_APPEND);
}
}
3 使用Prometheus + Grafana
// metrics端点
class DeadLetterMetrics
{
public function getMetrics()
{
$metrics = [];
// 死信队列大小
$metrics['dead_letter_queue_size'] = $this->getQueueSize('dead_letter');
// 死信处理失败次数
$metrics['dead_letter_process_failures'] = $this->getFailureCount();
// 死信重试次数分布
$metrics['dead_letter_retry_distribution'] = $this->getRetryDistribution();
return $metrics;
}
}
清理策略
1 自动清理脚本
<?php
class DeadLetterCleaner
{
private $maxRetries = 3;
private $retryDelay = 60; // 秒
// 清理死信消息
public function cleanDeadLetters()
{
try {
$deadMessages = $this->fetchDeadLetters();
foreach ($deadMessages as $message) {
$this->processDeadLetter($message);
}
} catch (Exception $e) {
$this->logError("清理死信失败: " . $e->getMessage());
}
}
// 处理单个死信
private function processDeadLetter($message)
{
$messageId = $message['message_id'];
$retryCount = $message['retry_count'] ?? 0;
// 1. 记录死信详情
$this->logDeadLetterDetail($message);
// 2. 判断重试次数
if ($retryCount < $this->maxRetries) {
// 重新投递到原始队列
$this->reQueueMessage($message);
} else {
// 3. 超过重试次数,归档或丢弃
$this->archiveDeadLetter($message);
// 或直接删除
$this->deleteDeadLetter($message);
}
// 4. 正式确认消费(ACK)
$this->acknowledgeMessage($messageId);
}
// 重新入队
private function reQueueMessage($message)
{
$channel = new AMQPChannel($this->connection);
$exchange = new AMQPExchange($channel);
$exchange->setName('original_exchange');
// 修改消息属性,增加重试计数
$properties = $message['properties'];
$properties['retry_count'] = ($message['retry_count'] ?? 0) + 1;
$properties['last_retry_time'] = time();
// 重新投递
$exchange->publish(
$message['body'],
$message['routing_key'],
AMQP_NOPARAM,
$properties
);
}
// 归档死信
private function archiveDeadLetter($message)
{
// 存储到数据库或文件系统
$archive = [
'message_id' => $message['message_id'],
'body' => $message['body'],
'headers' => $message['headers'],
'failed_at' => date('Y-m-d H:i:s'),
'retry_count' => $message['retry_count']
];
// 保存到MySQL
DB::table('dead_letter_archive')->insert($archive);
// 或保存到文件
file_put_contents(
"/var/dead_letters/{$message['message_id']}.json",
json_encode($archive)
);
}
// 批量清理过期死信
public function cleanExpiredDeadLetters($expireHours = 24)
{
$cutoffTime = time() - ($expireHours * 3600);
// 删除过期死信
DB::table('dead_letter_archive')
->where('created_at', '<', date('Y-m-d H:i:s', $cutoffTime))
->delete();
// 清理日志文件
$this->cleanExpiredLogs($expireHours);
}
}
2 定时任务清理
# crontab 配置 # 每分钟检查死信队列 * * * * * php /path/to/check_dead_letters.php # 每5分钟清理过期死信 */5 * * * * php /path/to/clean_dead_letters.php # 每小时汇报死信统计 0 * * * * php /path/to/report_dead_letters.php
3 基于Redis的清理
<?php
class RedisDeadLetterCleaner
{
private $redis;
public function __construct()
{
$this->redis = new Redis();
$this->redis->connect('127.0.0.1', 6379);
}
// 清理Redis死信
public function cleanRedisDeadLetters()
{
// 1. 获取所有死信key
$deadKeys = $this->redis->keys('dead_letter:*');
foreach ($deadKeys as $key) {
$deadData = $this->redis->hGetAll($key);
// 检查TTL
$ttl = $this->redis->ttl($key);
if ($ttl < 0) {
// 设置过期时间
$this->redis->expire($key, 86400); // 24小时
}
// 检查重试次数
$retryCount = $deadData['retry_count'] ?? 0;
if ($retryCount >= 3) {
// 移动到归档列表
$this->redis->rPush('dead_letter_archive', json_encode($deadData));
// 删除原始key
$this->redis->del($key);
}
}
}
// 获取死信统计
public function getDeadLetterStats()
{
return [
'total' => $this->redis->lLen('dead_letter_archive'),
'recent' => $this->redis->lRange('dead_letter_archive', 0, 10)
];
}
}
健康检查API
<?php
// health_check.php
class DeadLetterHealthCheck
{
public function check()
{
$status = [
'healthy' => true,
'dead_letter_queue_size' => 0,
'message' => 'OK'
];
try {
// 检查死信队列
$queueSize = $this->getDeadLetterQueueSize();
$status['dead_letter_queue_size'] = $queueSize;
if ($queueSize > 1000) {
$status['healthy'] = false;
$status['message'] = "死信队列积压: {$queueSize}条";
}
// 检查最近处理成功率
$successRate = $this->getRecentSuccessRate();
if ($successRate < 0.95) {
$status['healthy'] = false;
$status['message'] .= ", 处理成功率: {$successRate}%";
}
} catch (Exception $e) {
$status['healthy'] = false;
$status['message'] = "检查失败: " . $e->getMessage();
}
return $status;
}
private function getDeadLetterQueueSize()
{
// 实现具体队列大小获取逻辑
return 50;
}
private function getRecentSuccessRate()
{
// 计算最近5分钟的处理成功率
return 0.98;
}
}
// 返回健康检查结果
header('Content-Type: application/json');
$checker = new DeadLetterHealthCheck();
echo json_encode($checker->check());
最佳实践
1 死信处理流程
// 1. 定义重试策略
$retryPolicy = [
'max_retries' => 3,
'backoff_strategy' => 'exponential', // 指数退避
'initial_delay' => 5, // 秒
'multiplier' => 2
];
// 2. 死信分类处理
class DeadLetterClassifier
{
public function classify($message)
{
$errorCode = $message['headers']['error_code'] ?? 0;
return match(true) {
$errorCode >= 500 => 'system_error', // 系统错误,可重试
$errorCode >= 400 => 'client_error', // 客户端错误,不可重试
$message['headers']['timeout'] ?? false => 'timeout_error', // 超时
default => 'unknown'
};
}
}
// 3. 自动修复机制
class AutoRepairService
{
public function attemptAutoRepair($deadLetter)
{
$classifier = new DeadLetterClassifier();
$type = $classifier->classify($deadLetter);
return match($type) {
'system_error' => $this->retryWithBackoff($deadLetter),
'client_error' => $this->notifyAdmin($deadLetter),
'timeout_error' => $this->increaseTimeout($deadLetter),
default => $this->logAndArchive($deadLetter)
};
}
}
2 预警配置
// 预警阈值配置
$alertConfig = [
'critical_threshold' => 1000, // 严重
'warning_threshold' => 500, // 警告
'info_threshold' => 100, // 信息
// 监控周期
'check_interval' => 60, // 秒
// 通知渠道
'channels' => [
'email' => ['admin@example.com'],
'slack' => '#alerts',
'webhook' => 'https://hooks.example.com/dead_letter'
]
];
日志和审计
// 死信审计日志
class DeadLetterAudit
{
public function logDeadLetterEvent($event, $details)
{
$logEntry = [
'timestamp' => date('Y-m-d H:i:s'),
'event' => $event, // created, processed, archived, deleted
'details' => $details,
'source_ip' => $_SERVER['REMOTE_ADDR'] ?? 'unknown',
'user_agent' => $_SERVER['HTTP_USER_AGENT'] ?? 'unknown'
];
// 写入日志
$this->writeAuditLog($logEntry);
// 同步到监控系统
$this->syncToMonitor($logEntry);
}
}
一个完善的PHP死信监控清理系统应该:
- 实时监控:多种方式监控死信队列状态
- 自动清理:基于策略的自动清理机制
- 告警通知:不同的严重级别通知
- 历史记录:完整的审计日志
- 健康检查:API端点监控系统健康状态
- 容错处理:清理脚本自身的容错机制
根据实际项目规模和需求,可以选择合适的组件组合实现,小项目使用简单的脚本+定时任务即可,大型项目需要完整的监控体系。