PHP项目黑名单数据如何多节点实时同步

wen PHP项目 29

本文目录导读:

PHP项目黑名单数据如何多节点实时同步

  1. 基于Redis的Pub/Sub(推荐首选)
  2. 基于MySQL + 版本号轮询
  3. 基于消息队列(RabbitMQ/Kafka)
  4. 基于文件同步(简单场景)
  5. 最佳实践建议
  6. 推荐方案选择
  7. 注意事项

针对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 + 消息队列 灵活且可扩展

注意事项

  1. 网络分区处理:节点断连恢复后应自动全量同步
  2. 数据一致性:使用分布式锁或乐观锁避免并发冲突
  3. 监控告警:实时监控同步延迟和失败率
  4. 备份方案:定期将黑名单数据持久化到数据库

选择方案时,建议优先考虑Redis Pub/Sub方案,它在PHP生态中实现简单、性能优秀,适合大多数业务场景,如果对消息可靠性要求极高,可升级到RabbitMQ或Kafka。

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