本文目录导读:

- 使用官方 Kafka PHP 扩展(推荐)
- 使用 PHP Kafka 库(php-kafka)
- 通过 HTTP REST Proxy(简单方式)
- 异步日志处理方案(高性能)
- 配合 Monolog 使用
- 配置建议
- 注意事项
在 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})");
});
注意事项
- 错误处理:始终处理发送失败的情况,避免丢失日志
- 性能优化:使用异步发送、批量发送提高性能
- 消息压缩:对大量日志,启用压缩节省带宽
- 监控:添加监控指标,如发送成功率、延迟等
- 安全:生产环境配置 SASL/SSL 认证
选择哪种方式取决于你的需求:
- 需要高性能 → 使用 rdkafka 扩展
- 简单快速实现 → 使用 REST Proxy
- 需要日志框架集成 → 使用 Monolog
- 不想安装扩展 → 使用纯 PHP 库