PHP 发送日志到Kafka

wen PHP项目 3

本文目录导读:

PHP 发送日志到Kafka

  1. 使用官方 Kafka PHP 扩展(推荐)
  2. 使用 PHP Kafka 库(php-kafka)
  3. 通过 HTTP REST Proxy(简单方式)
  4. 异步日志处理方案(高性能)
  5. 配合 Monolog 使用
  6. 配置建议
  7. 注意事项

在 PHP 中发送日志到 Kafka 有多种方式,我来介绍几种常用的方法:

使用官方 Kafka PHP 扩展(推荐)

安装扩展

# 在 Ubuntu/Debian 上
pecl install rdkafka
# 添加扩展到 php.ini
echo "extension=rdkafka.so" >> /etc/php/8.2/cli/php.ini
# 需要先安装 librdkafka
sudo apt-get install librdkafka-dev

基础使用

<?php
class KafkaLogger {
    private $producer;
    private $topic;
    public function __construct($brokers, $topicName) {
        $conf = new RdKafka\Conf();
        // 设置生产者配置
        $conf->set('metadata.broker.list', $brokers);
        $conf->set('socket.timeout.ms', 100);
        $conf->set('message.timeout.ms', 5000);
        // 设置回调
        $conf->setDrMsgCb(function ($kafka, $message) {
            if ($message->err) {
                error_log("Kafka 发送失败: " . $message->errstr());
            } else {
                error_log("Kafka 发送成功: " . $message->offset);
            }
        });
        $this->producer = new RdKafka\Producer($conf);
        $this->topic = $this->producer->newTopic($topicName);
    }
    public function send(array $logData) {
        try {
            $message = json_encode($logData);
            // 发送消息(分区0)
            $this->topic->produce(
                RD_KAFKA_PARTITION_UA, // 自动分配分区
                0,                     // 消息key
                $message
            );
            // 刷新到Kafka
            $this->producer->flush(1000);
            // 轮询回调
            $this->producer->poll(0);
            return true;
        } catch (Exception $e) {
            error_log("Kafka 发送异常: " . $e->getMessage());
            return false;
        }
    }
    public function close() {
        $this->producer->flush(5000);
    }
}
// 使用示例
$logger = new KafkaLogger('localhost:9092,localhost:9093', 'app_logs');
$logData = [
    'timestamp' => date('Y-m-d H:i:s'),
    'level' => 'INFO',
    'message' => '用户登录成功',
    'user_id' => 12345,
    'ip' => $_SERVER['REMOTE_ADDR']
];
$logger->send($logData);
$logger->close();

使用 PHP Kafka 库(php-kafka)

安装

composer require weiboad/kafka-php

使用示例

<?php
require_once 'vendor/autoload.php';
use Kafka\Producer;
use Kafka\ProducerConfig;
class KafkaLogger {
    private $producer;
    public function __construct($brokers) {
        $config = ProducerConfig::getInstance();
        // 配置Kafka连接
        $config->setMetadataBrokerList($brokers);
        $config->setRequiredAck(1);          // 等待leader确认
        $config->setIsAsyn(false);            // 同步发送
        $this->producer = new Producer();
    }
    public function send($topic, array $logData) {
        try {
            $message = [
                'topic' => $topic,
                'key' => null,
                'value' => json_encode($logData),
                'timestamp' => time() * 1000
            ];
            $this->producer->send([$message]);
            return true;
        } catch (Exception $e) {
            error_log("Kafka 发送异常: " . $e->getMessage());
            return false;
        }
    }
}
// 使用示例
$logger = new KafkaLogger('localhost:9092');
$logger->send('app_logs', [
    'timestamp' => date('Y-m-d H:i:s'),
    'level' => 'ERROR',
    'message' => '数据库连接失败'
]);

通过 HTTP REST Proxy(简单方式)

使用 Kafka REST Proxy

<?php
class KafkaRestLogger {
    private $restEndpoint;
    private $topic;
    public function __construct($restEndpoint, $topic) {
        $this->restEndpoint = $restEndpoint;
        $this->topic = $topic;
    }
    public function send(array $logData) {
        $ch = curl_init();
        $payload = [
            'records' => [
                [
                    'value' => $logData
                ]
            ]
        ];
        curl_setopt($ch, CURLOPT_URL, $this->restEndpoint . '/topics/' . $this->topic);
        curl_setopt($ch, CURLOPT_POST, true);
        curl_setopt($ch, CURLOPT_HTTPHEADER, [
            'Content-Type: application/vnd.kafka.json.v2+json'
        ]);
        curl_setopt($ch, CURLOPT_POSTFIELDS, json_encode($payload));
        curl_setopt($ch, CURLOPT_RETURNTRANSFER, true);
        $response = curl_exec($ch);
        $httpCode = curl_getinfo($ch, CURLINFO_HTTP_CODE);
        curl_close($ch);
        return $httpCode == 200;
    }
}
// 使用示例
$logger = new KafkaRestLogger('http://localhost:8082', 'app_logs');
$logger->send([
    'timestamp' => time(),
    'level' => 'WARNING',
    'message' => '磁盘空间不足'
]);

异步日志处理方案(高性能)

<?php
class AsyncKafkaLogger {
    private $kafkaProducer;
    private $processQueue = [];
    private $maxQueueSize = 100;
    private $flushInterval = 5; // 秒
    public function __construct($brokers, $topic) {
        $conf = new RdKafka\Conf();
        $conf->set('metadata.broker.list', $brokers);
        // 启用异步和批量
        $conf->set('queue.buffering.max.messages', 10000);
        $conf->set('batch.num.messages', 100);
        $conf->set('compression.type', 'snappy');
        $this->producer = new RdKafka\Producer($conf);
        $this->topic = $this->producer->newTopic($topic);
        // 设置定时刷新
        register_shutdown_function([$this, 'flushAll']);
    }
    public function log($level, $message, array $context = []) {
        $logEntry = [
            'timestamp' => date('c', time()),
            'level' => $level,
            'message' => $message,
            'context' => $context,
            'host' => gethostname(),
            'process_id' => getmypid()
        ];
        $this->topic->produce(
            RD_KAFKA_PARTITION_UA,
            0,
            json_encode($logEntry)
        );
        // 每100条消息刷新一次
        if (count($this->processQueue) >= $this->maxQueueSize) {
            $this->flushAll();
        }
    }
    public function flushAll() {
        $this->producer->flush(1000);
        $this->producer->poll(0);
    }
    public function logInfo($message, array $context = []) {
        $this->log('INFO', $message, $context);
    }
    public function logError($message, array $context = []) {
        $this->log('ERROR', $message, $context);
    }
    public function logWarning($message, array $context = []) {
        $this->log('WARNING', $message, $context);
    }
}
// 使用示例
$logger = new AsyncKafkaLogger('localhost:9092', 'app_logs');
// 业务代码中记录日志
$logger->logInfo('用户登录', ['user_id' => 123]);
$logger->logError('支付失败', ['order_id' => 'ORD-2024-123']);
// 在脚本结束时自动flush

配合 Monolog 使用

<?php
require_once 'vendor/autoload.php';
use Monolog\Logger;
use Monolog\Handler\KafkaHandler;
// 自定义 Kafka Handler
class KafkaMonologHandler extends AbstractProcessingHandler {
    private $producer;
    private $topic;
    public function __construct($brokers, $topic, $level = Logger::DEBUG) {
        parent::__construct($level);
        $conf = new RdKafka\Conf();
        $conf->set('metadata.broker.list', $brokers);
        $this->producer = new RdKafka\Producer($conf);
        $this->topic = $this->producer->newTopic($topic);
    }
    protected function write(array $record): void {
        $message = json_encode($record);
        $this->topic->produce(RD_KAFKA_PARTITION_UA, 0, $message);
        $this->producer->flush(100);
    }
}
// 创建 Monolog logger
$logger = new Logger('app');
$logger->pushHandler(new KafkaMonologHandler('localhost:9092', 'monolog_logs'));
// 使用
$logger->info('Something happened', ['extra_data' => 'value']);

配置建议

生产环境最佳实践

$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', 'kafka1:9092,kafka2:9092,kafka3:9092');
$conf->set('client.id', 'php-log-producer');
// 可靠性设置
$conf->set('request.required.acks', 1);        // 至少收到leader确认
$conf->set('message.send.max.retries', 3);      // 重试次数
$conf->set('retry.backoff.ms', 100);            // 重试间隔
// 批量设置
$conf->set('batch.num.messages', 100);          // 批量消息数
$conf->set('queue.buffering.max.messages', 100000); // 队列缓冲
// 压缩
$conf->set('compression.type', 'snappy');       // 压缩类型
// 错误处理
$conf->setErrorCb(function ($kafka, $err, $reason) {
    error_log("Kafka 错误: {$reason} (代码: {$err})");
});

注意事项

  1. 错误处理:始终处理发送失败的情况,避免丢失日志
  2. 性能优化:使用异步发送、批量发送提高性能
  3. 消息压缩:对大量日志,启用压缩节省带宽
  4. 监控:添加监控指标,如发送成功率、延迟等
  5. 安全:生产环境配置 SASL/SSL 认证

选择哪种方式取决于你的需求:

  • 需要高性能 → 使用 rdkafka 扩展
  • 简单快速实现 → 使用 REST Proxy
  • 需要日志框架集成 → 使用 Monolog
  • 不想安装扩展 → 使用纯 PHP 库

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