
前言用 PHP 消费 Kafka第一次上线常见的三种症状是消息重复消费同一条订单处理了三遍、消息静默丢失offset 提交了但业务逻辑抛了异常、以及消费者反复被踢出组日志里刷Group coordinator is rebalancing消费速率趋近于零。这三种症状对应的是同一件事的三个侧面提交位点offset的时机。PHP 不像 Java 那样有常驻线程模型一次请求结束进程就没了所以「消费者」在 PHP 里必须是一个常驻 CLI 进程。这个前提决定了所有配置项的选择逻辑消费者要长生命周期、要手动控制提交、要能优雅退出、要控制单批处理时长否则协调者会认为它已经死了。本文用 PHP 的rdkafka扩展php-rdkafka底层是 librdkafka写一个完整可运行的消费者手动提交位点、正常处理空轮询、支持优雅退出。示例代码需要PHP 8.1 及以上我们在 PHP 8.5 上运行——这个话题和 PHP 版本关系不大真正需要对齐的是扩展版本与 librdkafka 的版本。一、先理解消费模型里的三个角色Kafka 的消费语义由三个概念共同决定配置项选错了语义就跟着错概念含义决定了什么消费者组consumer group一组消费者共享一个group.id同组内每条消息只被一个消费者处理分区partition主题的并行单元一个分区同一时刻只属于组内一个消费者分区数就是并行度上限位点offset每个分区上「已消费到哪」的标记提交时机决定了重复还是丢失由这三者推出的几条硬结论消费者数量超过分区数多出来的消费者会一直空闲。想提高并行度先加分区。enable.auto.commit打开时位点是按时间自动提交的业务还没跑完位点可能已经提交了。一旦进程崩掉那条消息就再也回不来了——这叫 at-most-once至多一次。手动提交 先处理再提交得到的是 at-least-once至少一次代价是崩溃后可能重复处理一条消息所以业务逻辑必须幂等。绝大多数业务场景要的就是「至少一次 幂等」而不是去追求「恰好一次」。把幂等键订单号、消息 key落在数据库唯一索引上比配置一长串事务参数可靠得多。二、装扩展并确认版本匹配php-rdkafka 是 PECL 扩展依赖系统里的 librdkafka# Debian / Ubuntu先装 librdkafka 开发包 sudo apt-get install -y librdkafka-dev # 再装 PHP 扩展 sudo pecl install rdkafka # 在 php.ini 里启用 # extensionrdkafka确认装好了php -m | grep rdkafka php -r var_dump(extension_loaded(rdkafka), RdKafka\LIBRDKAFKA_VERSION);这一步容易出问题的地方扩展是编译型扩展必须为每个 PHP 版本单独编译。升级到 PHP 8.5 之后原来为 8.2 编译的rdkafka.so直接加载失败php -m里看不到它——症状是「升级后队列消费脚本启动就白屏」。所以升级 PHP 主版本时要同步确认这些 PECL 扩展有没有对应的兼容版本。如果不想依赖扩展社区也有纯 PHP 实现的 Kafka 客户端但功能和性能取舍要自己评估。三、配置项只关心这几个RdKafka\Conf是一层键值配置键名就是 librdkafka 的配置名。下面这些是必须显式设置的配置键建议值为什么metadata.broker.list你的 broker 地址不填连不上group.id业务语义化的名字决定「谁和谁是一组」auto.offset.resetearliest或latest无位点时从哪开始默认是latestlibrdkafka 里叫largest不设会「丢历史消息」enable.auto.commitfalse手动提交位点才可控session.timeout.ms按 broker 版本建议值超时会话被判定死亡触发 rebalancemax.poll.interval.ms大于单批最长处理时间超了就被踢出组导致重复消费关于后两项具体默认值随 librdkafka 版本变化以你所装版本的官方文档为准要点是「max.poll.interval.ms必须大于你一轮consume()到下一次consume()之间的最长耗时」否则消费者会不断被踢出组。实战完整可运行的消费者把下面这段保存成consume.php。它需要三个前置条件PHP 8.1、已加载rdkafka扩展、可连通的 Kafka broker。?php declare(strict_types1); // consume.php —— 需要 PHP 8.1 与 ext-rdkafka // 用法: php consume.php topic [消费者组名] if (!extension_loaded(rdkafka) || !class_exists(\RdKafka\KafkaConsumer::class)) { exit(需要 ext-rdkafka请先 pecl install rdkafka 并在 php.ini 里启用\n); } $topic $argv[1] ?? ; $groupId $argv[2] ?? php-demo-group; if ($topic ) { exit(用法: php consume.php topic [消费者组名]\n); } // 优雅退出收到 SIGTERM / SIGINT 时把标志置位主循环跑完当前消息再退出 $running true; if (function_exists(pcntl_signal)) { pcntl_async_signals(true); pcntl_signal(SIGTERM, static function () use ($running): void { $running false; }); pcntl_signal(SIGINT, static function () use ($running): void { $running false; }); } $conf new \RdKafka\Conf(); $conf-set(metadata.broker.list, 127.0.0.1:9092); $conf-set(group.id, $groupId); $conf-set(auto.offset.reset, earliest); // 没有位点时从头消费 $conf-set(enable.auto.commit, false); // 关键手动提交 // rebalance 回调分区被收走时要把「正在处理的进度」落盘否则会重复消费 $conf-setRebalanceCb( static function (\RdKafka\KafkaConsumer $consumer, int $err, ?array $partitions null): void { switch ($err) { case RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS: $consumer-assign($partitions); break; case RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS: $consumer-assign(null); break; default: $consumer-assign(null); break; } } ); $consumer new \RdKafka\KafkaConsumer($conf); $consumer-subscribe([$topic]); echo 已订阅 {$topic}group{$groupId}等待消息…\n; /** * 业务处理。真实项目里这里应当只做「幂等」的落库/调用 * 用消息 key 或业务主键做唯一索引重复执行也不会产生副作用。 */ function handle(\RdKafka\Message $message): void { printf( [%s] partition%d offset%d key%s payload%s\n, $message-topic_name, $message-partition, $message-offset, var_export($message-key, true), substr((string) $message-payload, 0, 200) ); } while ($running) { // consume() 的超时单位是毫秒超时返回一条 err 为 RD_KAFKA_RESP_ERR__TIMED_OUT 的「空消息」 $message $consumer-consume(1000); switch ($message-err) { case RD_KAFKA_RESP_ERR_NO_ERROR: try { handle($message); } catch (Throwable $e) { // 处理失败不提交位点让这条消息下次重新投递 fwrite(STDERR, 处理失败跳过提交{$e-getMessage()}\n); continue 2; } // 只有处理成功才提交位点同步提交返回后位点已生效 $consumer-commit($message); break; case RD_KAFKA_RESP_ERR__PARTITION_EOF: // 追上了分区末尾不是错误继续轮询即可等到新消息 break; case RD_KAFKA_RESP_ERR__TIMED_OUT: // 这段时间没有消息正常现象 break; default: fwrite(STDERR, 消费出错{$message-errstr()}code{$message-err}\n); // 不要直接 break 整个循环短暂的协调者切换会自愈 usleep(500_000); break; } } $consumer-close(); echo 已停止消费并释放分区\n;运行php consume.php orders php-demo-group输出形如已订阅 ordersgroupphp-demo-group等待消息… [orders] partition0 offset17 keyorder-1001 payload{id:1001,amount:19.9} [orders] partition2 offset42 keyorder-1002 payload{id:1002,amount:5.0}想验证「处理失败不丢消息」把handle()改成对特定 payload 抛异常再重跑你会看到同一条消息被再次投递——这正是 at-least-once 的预期行为。常见坑点1. 开着自动提交却指望业务失败能重试❌ 用默认的enable.auto.committrue业务逻辑抛异常后进程退出位点早就被后台自动提交了那条消息永远不会再来。 ✅ 设enable.auto.commitfalse处理成功后再显式commit()业务逻辑做成幂等接受「可能重复」。2. 把RD_KAFKA_RESP_ERR__TIMED_OUT当成故障❌if ($message-err ! RD_KAFKA_RESP_ERR_NO_ERROR) { exit(1); }—— 队列空闲一会儿消费者进程就自己退出了。 ✅ 把__TIMED_OUT和__PARTITION_EOF都当成「正常空轮询」继续循环只有真正的错误码才记日志并退避重试。3. 每处理一条消息就新建一次消费者❌ 在循环体内new \RdKafka\KafkaConsumer($conf)再subscribe()每次都会触发一次组内 rebalance消费速率断崖式下跌。 ✅ 消费者对象在进程里只建一次长期存活进程重启是运维动作不是每条消息的动作。4.max.poll.interval.ms小于单批处理时间❌ 一条消息要处理 10 分钟比如调用外部慢接口而轮询间隔上限是 5 分钟协调者判定消费者已死把它踢出组分区分给别人 → 同一条消息被两处处理。 ✅ 让单批处理时间远小于max.poll.interval.ms长任务拆成小步骤或投递到别的队列不要卡在消费循环里。5. rebalance 回调里不做任何清理❌ 只给RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS写assign()撤销分支随便处理正在处理的 offset 没保存。 ✅ 撤销时先把「已处理到哪」记下来或直接commit再assign(null)交出分区避免重复处理。6. 消费者数量超过分区数❌ 订单主题只有 3 个分区却起了 10 个消费者进程7 个进程长期空转还占着连接和内存。 ✅ 并行度上限就是分区数先扩分区注意扩分区会打乱 key 到分区的映射再考虑加消费者。7. 忘了close()❌ 进程直接exit消费者没有离开组协调者要等到会话超时才触发 rebalance这段时间同组其他消费者都在等。 ✅ 退出前调用close()它会提交最后的位点并主动离开组配合信号处理做优雅退出。8. 用latest却以为能消费到历史数据❌ 新起一个消费者组auto.offset.reset保持默认结果「一条历史消息都没收到」以为代码有问题。 ✅ 明确这个配置的语义earliest从头latest只消费启动之后的新消息。调试阶段用earliest生产按业务需要选。总结环节建议避免的问题消费者形态常驻 CLI 进程每请求重建导致的 rebalance位点提交enable.auto.commitfalse 处理成功再提交消息静默丢失业务逻辑幂等业务键唯一索引重复消费造成脏数据空轮询忽略__TIMED_OUT/__PARTITION_EOF空闲即退出并行度消费者数 ≤ 分区数进程空转退出信号处理 close()长时间 rebalance 空窗PHP 消费 Kafka 的难点不在语法而在给自己立下几条纪律进程常驻、位点手动、业务幂等、退出优雅。把这四条写进代码模板里剩下的配置项就只是填空反过来只要「先提交后处理」这一条错了无论怎么调参数消息迟早会丢。