You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用nmerd/kafka-php异步发送消息异常:消息丢失但成功日志有记录

问题:nmerd/kafka-php异步发送消息丢失,日志异常

异常现象

  • 发送2条消息仅能收到1条,发送3条仅能收到2条,但成功日志显示所有消息已发送成功
  • 日志表现异常:
    • 发送2条消息时,仅出现1个唯一key,对应Producer created日志打印2次
    • 发送3条消息时,仅出现2个唯一key,其中一个key的日志打印1次,另一个打印2次

相关代码

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);
    }
}

问题分析

  1. 配置模式矛盾:异步发送方法调用的initializeConfig强制设置setIsAsyn(false)(同步模式),但代码却使用了异步生产者的构造逻辑(传入消息生成闭包),导致发送机制混乱,消息被错误合并或丢失。
  2. 全局配置重复覆盖:每次调用发送方法都会重置全局ProducerConfig实例,并发场景下后一次配置会覆盖前一次,导致前一次的生产者实例行为异常,出现日志key重复、消息丢失的情况。
  3. 批量配置误用: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 13:37:03