本文目录导读:

针对PHP项目的黑名单数据多节点实时同步,以下是几种主流且实用的解决方案,按推荐优先级排序:
基于Redis的Pub/Sub(推荐首选)
架构设计
// 黑名单发布者(修改黑名单时触发)
class BlacklistPublisher {
private $redis;
public function __construct($redis) {
$this->redis = $redis;
}
public function publishUpdate($action, $data) {
$message = json_encode([
'action' => $action, // 'add', 'remove', 'clear'
'data' => $data,
'timestamp' => time(),
'node_id' => gethostname()
]);
// 发布到黑名单频道
$this->redis->publish('blacklist:updates', $message);
// 同时更新本地缓存
$this->updateLocalCache($action, $data);
}
private function updateLocalCache($action, $data) {
// 更新本地Redis或内存缓存
}
}
// 黑名单订阅者(每个节点启动时注册)
class BlacklistSubscriber {
private $redis;
private $localCache;
public function __construct($redis) {
$this->redis = $redis;
$this->localCache = new BlacklistLocalCache();
}
public function subscribe() {
// 创建订阅循环(建议在独立进程中运行)
$this->redis->subscribe(['blacklist:updates'], function($redis, $channel, $message) {
$this->handleUpdate($message);
});
}
private function handleUpdate($message) {
$data = json_decode($message, true);
switch($data['action']) {
case 'add':
$this->localCache->add($data['data']);
break;
case 'remove':
$this->localCache->remove($data['data']);
break;
case 'clear':
$this->localCache->clear();
break;
}
// 记录同步日志
$this->logSync($data);
}
}
启动订阅进程
# 使用Supervisor管理订阅进程 [program:blacklist-subscriber] command=php /path/to/blacklist_subscriber.php process_name=%(program_name)s_%(process_num)02d numprocs=1 autostart=true autorestart=true user=www-data
基于MySQL + 版本号轮询
数据库设计
CREATE TABLE blacklist_sync (
id INT PRIMARY KEY AUTO_INCREMENT,
data_type VARCHAR(50) NOT NULL, -- IP, IP段, 用户ID等
data_value VARCHAR(255) NOT NULL,
action ENUM('add', 'remove') NOT NULL,
version BIGINT NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
node_id VARCHAR(100),
INDEX idx_version (version)
);
-- 黑名单主表
CREATE TABLE blacklist_items (
id INT PRIMARY KEY AUTO_INCREMENT,
data_type VARCHAR(50) NOT NULL,
data_value VARCHAR(255) NOT NULL,
expires_at TIMESTAMP NULL,
status TINYINT DEFAULT 1,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_type_value (data_type, data_value)
);
PHP实现
class BlacklistSyncManager {
private $db;
private $localVersion = 0;
public function __construct() {
$this->db = Database::getConnection();
// 初始化时获取本地版本号
$this->localVersion = $this->getLocalVersion();
}
public function addBlacklist($type, $value) {
try {
$this->db->beginTransaction();
// 插入主表
$stmt = $this->db->prepare(
"INSERT INTO blacklist_items (data_type, data_value, status)
VALUES (?, ?, 1)
ON DUPLICATE KEY UPDATE status = 1"
);
$stmt->execute([$type, $value]);
// 计算新版本号
$newVersion = $this->getNextVersion();
// 插入同步记录
$stmt = $this->db->prepare(
"INSERT INTO blacklist_sync (data_type, data_value, action, version, node_id)
VALUES (?, ?, 'add', ?, ?)"
);
$stmt->execute([$type, $value, $newVersion, gethostname()]);
$this->db->commit();
// 更新本地版本号
$this->localVersion = $newVersion;
} catch (Exception $e) {
$this->db->rollBack();
throw $e;
}
}
public function syncFromRemote() {
// 查询自上次同步以来的变更
$stmt = $this->db->prepare(
"SELECT * FROM blacklist_sync WHERE version > ? ORDER BY version ASC"
);
$stmt->execute([$this->localVersion]);
$changes = $stmt->fetchAll(PDO::FETCH_ASSOC);
foreach ($changes as $change) {
// 更新本地缓存
$this->applyChange($change);
$this->localVersion = $change['version'];
}
return count($changes);
}
private function applyChange($change) {
// 根据action更新本地缓存
if ($change['action'] == 'add') {
// 添加到本地缓存(Redis/内存/文件)
} else {
// 从本地缓存移除
}
}
private function getNextVersion() {
// 使用Redis INCR或数据库自增获取版本号
return $this->db->query("SELECT UNIX_TIMESTAMP(NOW()) * 1000 + RAND() * 1000")->fetchColumn();
}
}
基于消息队列(RabbitMQ/Kafka)
RabbitMQ实现
// 生产者服务
class BlacklistProducer {
private $channel;
public function __construct() {
$connection = new AMQPStreamConnection('localhost', 5672, 'user', 'pass');
$this->channel = $connection->channel();
$this->channel->exchange_declare('blacklist_exchange', 'fanout', false, true, false);
}
public function publish($action, $data) {
$message = new AMQPMessage(json_encode([
'action' => $action,
'data' => $data,
'node_id' => gethostname()
]));
$this->channel->basic_publish($message, 'blacklist_exchange');
}
}
// 消费者服务(每个节点运行)
class BlacklistConsumer {
private $channel;
private $localCache;
public function __construct() {
$connection = new AMQPStreamConnection('localhost', 5672, 'user', 'pass');
$this->channel = $connection->channel();
$this->channel->exchange_declare('blacklist_exchange', 'fanout', false, true, false);
// 每个节点创建独立队列
list($queueName) = $this->channel->queue_declare("", false, true, true, true);
$this->channel->queue_bind($queueName, 'blacklist_exchange');
$this->localCache = new BlacklistLocalCache();
}
public function consume() {
$this->channel->basic_consume($queueName, '', false, true, false, false, function($msg) {
$this->handleMessage($msg->body);
});
while (count($this->channel->callbacks)) {
$this->channel->wait();
}
}
private function handleMessage($body) {
$data = json_decode($body, true);
$this->localCache->sync($data);
}
}
基于文件同步(简单场景)
使用NFS/Samba共享
class SharedFileBlacklist {
private $sharedFile = '/mnt/shared/blacklist.json';
private $localCacheFile = '/tmp/blacklist_cache.json';
private $lockFile = '/mnt/shared/blacklist.lock';
public function add($type, $value) {
$lock = $this->acquireLock();
$blacklist = $this->readSharedFile();
$blacklist[$type][] = $value;
$this->writeSharedFile($blacklist);
$this->releaseLock($lock);
// 通知其他节点(可选)
$this->notifyNodes('add', $type, $value);
}
private function acquireLock() {
$fp = fopen($this->lockFile, 'w');
if (flock($fp, LOCK_EX)) {
return $fp;
}
throw new Exception('Could not acquire lock');
}
private function readSharedFile() {
if (file_exists($this->sharedFile)) {
return json_decode(file_get_contents($this->sharedFile), true) ?: [];
}
return [];
}
private function writeSharedFile($data) {
file_put_contents($this->sharedFile, json_encode($data), LOCK_EX);
}
}
最佳实践建议
缓存策略优化
class BlacklistCache {
private $redis;
private $memoryCache = [];
// 多级缓存读取
public function isBlacklisted($type, $value) {
// 1. 检查内存缓存(最快)
if (isset($this->memoryCache[$type][$value])) {
return $this->memoryCache[$type][$value];
}
// 2. 检查Redis
$key = "blacklist:{$type}";
if ($this->redis->sIsMember($key, $value)) {
$this->memoryCache[$type][$value] = true;
return true;
}
return false;
}
// 批量预加载
public function preloadBlacklist($type) {
$key = "blacklist:{$type}";
$members = $this->redis->sMembers($key);
foreach ($members as $member) {
$this->memoryCache[$type][$member] = true;
}
}
}
健康检查与故障恢复
class BlacklistHealthCheck {
public function checkSyncStatus() {
$status = [
'local_version' => $this->getLocalVersion(),
'remote_version' => $this->getRemoteVersion(),
'is_synced' => true,
'last_sync_time' => $this->getLastSyncTime(),
'queue_size' => $this->getPendingQueueSize()
];
$status['is_synced'] = ($status['local_version'] >= $status['remote_version']);
return $status;
}
public function forceResync() {
// 全量重新同步
$this->clearLocalCache();
$this->fullSyncFromDB();
}
}
推荐方案选择
| 场景 | 推荐方案 | 理由 |
|---|---|---|
| 中小规模(<10节点) | Redis Pub/Sub | 简单高效,延迟低 |
| 大规模(>10节点) | RabbitMQ/Kafka | 消息持久化,支持分布式 |
| 已有MySQL集群 | 版本号轮询 | 无需引入新组件 |
| 简单需求 | 文件共享 | 实现简单,但性能有限 |
| 微服务架构 | Redis + 消息队列 | 灵活且可扩展 |
注意事项
- 网络分区处理:节点断连恢复后应自动全量同步
- 数据一致性:使用分布式锁或乐观锁避免并发冲突
- 监控告警:实时监控同步延迟和失败率
- 备份方案:定期将黑名单数据持久化到数据库
选择方案时,建议优先考虑Redis Pub/Sub方案,它在PHP生态中实现简单、性能优秀,适合大多数业务场景,如果对消息可靠性要求极高,可升级到RabbitMQ或Kafka。