使用nmerd/kafka-php异步发送消息异常:消息丢失但成功日志有记录
问题:nmerd/kafka-php异步发送消息丢失,日志异常
异常现象
- 发送2条消息仅能收到1条,发送3条仅能收到2条,但成功日志显示所有消息已发送成功
- 日志表现异常:
- 发送2条消息时,仅出现1个唯一key,对应
Producer created日志打印2次 - 发送3条消息时,仅出现2个唯一key,其中一个key的日志打印1次,另一个打印2次
- 发送2条消息时,仅出现1个唯一key,对应
相关代码
class KafkaPublisher { const LOG_NAMESPACE = 'kafka'; public static function publishAsynchronously($servers, $message, $key, $topic) { try { self::initializeConfig($servers); $producer = new Producer( function () use ($message, $key, $topic) { // 发送2条消息时,仅1个唯一key,日志打印2次 // 发送3条消息时,仅2个唯一key,一个打印1次,另一个打印2次 Logger::quickInfo( 'Producer created', [ 'key' => $key, 'message' => $message, ], self::LOG_NAMESPACE ); return [ [ 'topic' => $topic, 'value' => $message, 'timestamp' => time(), 'key' => json_encode($key), ], ]; } ); $producer->success( function ($result) use ($message) { // 发送2条显示2条唯一消息,发送3条显示3条唯一消息 Logger::quickInfo( 'Successfully posted to Kafka', [ 'result' => $result, 'message_sent' => $message, ], self::LOG_NAMESPACE ); } ); $producer->error( function ($errorCode) use ($message) { Logger::quickError( 'Unsuccessfully posted to Kafka', [ 'errorCode' => $errorCode, 'message_sent' => $message, ], self::LOG_NAMESPACE ); } ); $producer->send(true); } catch (Exception $e) { Logger::quickError( 'Unsuccessfully connected to Kafka', [ 'errorCode' => $e->getCode(), 'errorMessage' => $e->getMessage(), 'message_sent' => $message, ], self::LOG_NAMESPACE ); } } public static function publishSynchronously($servers, $message, $key, $topic) { try { self::initializeConfig($servers); $producer = new Producer(); $producer->send([ [ 'topic' => $topic, 'value' => $message, 'timestamp' => time(), 'key' => json_encode($key), ], ] ); } catch (Exception $e) { Logger::quickError( 'Unsuccessfully connected to Kafka', [ 'errorCode' => $e->getCode(), 'errorMessage' => $e->getMessage(), 'message_sent' => $message, ], self::LOG_NAMESPACE ); } } protected static function initializeConfig($servers) { $config = ProducerConfig::getInstance(); $config->setMetadataRefreshIntervalMs(10000); $config->setMetadataBrokerList(implode(',', $servers)); $config->setBrokerVersion('1.0.0'); $config->setRequiredAck(1); $config->setIsAsyn(false); $config->setProduceInterval(500); } }
问题分析
- 配置模式矛盾:异步发送方法调用的
initializeConfig强制设置setIsAsyn(false)(同步模式),但代码却使用了异步生产者的构造逻辑(传入消息生成闭包),导致发送机制混乱,消息被错误合并或丢失。 - 全局配置重复覆盖:每次调用发送方法都会重置全局
ProducerConfig实例,并发场景下后一次配置会覆盖前一次,导致前一次的生产者实例行为异常,出现日志key重复、消息丢失的情况。 - 批量配置误用:
setProduceInterval(500)是异步批量发送的间隔配置,但同步模式下该配置无效,混合使用会触发逻辑冲突。
解决方案
1. 拆分同步/异步配置
修改配置初始化方法,支持根据发送模式切换配置:
protected static function initializeConfig($servers, $isAsync = false) { $config = ProducerConfig::getInstance(); $config->setMetadataRefreshIntervalMs(10000); $config->setMetadataBrokerList(implode(',', $servers)); $config->setBrokerVersion('1.0.0'); $config->setRequiredAck(1); // 根据模式设置异步开关 $config->setIsAsyn($isAsync); // 仅异步模式启用批量发送间隔 if ($isAsync) { $config->setProduceInterval(500); } }
2. 修正异步发送方法的配置调用
在异步发送方法中传入true启用异步模式:
public static function publishAsynchronously($servers, $message, $key, $topic) { try { // 初始化异步模式配置 self::initializeConfig($servers, true); $producer = new Producer( function () use ($message, $key, $topic) { Logger::quickInfo( 'Producer created', [ 'key' => $key, 'message' => $message, ], self::LOG_NAMESPACE ); return [ [ 'topic' => $topic, 'value' => $message, 'timestamp' => time(), 'key' => json_encode($key), ], ]; } ); // 原有回调代码保持不变 $producer->success(function ($result) use ($message) { Logger::quickInfo( 'Successfully posted to Kafka', [ 'result' => $result, 'message_sent' => $message, ], self::LOG_NAMESPACE ); }); $producer->error(function ($errorCode) use ($message) { Logger::quickError( 'Unsuccessfully posted to Kafka', [ 'errorCode' => $errorCode, 'message_sent' => $message, ], self::LOG_NAMESPACE ); }); $producer->send(true); } catch (Exception $e) { Logger::quickError( 'Unsuccessfully connected to Kafka', [ 'errorCode' => $e->getCode(), 'errorMessage' => $e->getMessage(), 'message_sent' => $message, ], self::LOG_NAMESPACE ); } }
3. 避免全局配置频繁重置(可选优化)
高频调用场景下,建议将配置初始化放到类的静态构造中,避免每次调用都重置全局配置:
private static $configInitialized = false; protected static function initializeConfig($servers) { if (self::$configInitialized) { return; } $config = ProducerConfig::getInstance(); $config->setMetadataRefreshIntervalMs(10000); $config->setMetadataBrokerList(implode(',', $servers)); $config->setBrokerVersion('1.0.0'); $config->setRequiredAck(1); self::$configInitialized = true; } // 异步方法中单独设置异步相关配置 public static function publishAsynchronously($servers, $message, $key, $topic) { try { self::initializeConfig($servers); $config = ProducerConfig::getInstance(); $config->setIsAsyn(true); $config->setProduceInterval(500); // 原有生产者代码保持不变 $producer = new Producer(function () use ($message, $key, $topic) { Logger::quickInfo( 'Producer created', [ 'key' => $key, 'message' => $message, ], self::LOG_NAMESPACE ); return [ [ 'topic' => $topic, 'value' => $message, 'timestamp' => time(), 'key' => json_encode($key), ], ]; }); // 回调及发送代码保持不变 $producer->success(...); $producer->error(...); $producer->send(true); } catch (Exception $e) { // 异常处理保持不变 } }
4. 验证消息真实性
在消费者端打印消息的key和offset,确认是否真的丢失;同时检查Kafka集群的acks配置是否与客户端setRequiredAck(1)匹配,确保消息被Broker持久化。
内容的提问来源于stack exchange,提问作者Nick
相关产品推荐
相关产品推荐

