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

使用Librdkafka向Kafka服务器发送消息时出现内存泄漏问题求助

Librdkafka向Kafka服务器发送消息时出现内存泄漏问题求助

大家好,我正在开发一个C语言应用,功能是将传感器数据封装成JSON字符串后发送到Kafka服务器。目前程序运行一段时间后会被Ubuntu系统强制终止,我怀疑是内存泄漏导致的,用Valgrind检测后得到了一些线索,但不确定问题出在我对librdkafka的使用上,还是库本身的bug,想请各位帮忙排查一下。

我的Kafka生产者初始化代码:

int setupKafkaProducer(struct KafkaParameters *kafkaParameters, struct ClientOPCEndpointInfo* *clientInfos, int clientInfosLength, bool runTest)
{
    logInfo("START - Setting up kafka producer", true);

    conf = rd_kafka_conf_new();
    rd_kafka_conf_set_dr_msg_cb(conf, dr_msg_cb);
    rd_kafka_conf_res_t res = RD_KAFKA_CONF_OK;

    // setting up parameters ...

    if (res != RD_KAFKA_CONF_OK)
    {
        g_error("Failed to setup kafka config: %s", errstr);
        logError("Failed to setup kafka config", true);
        return 1;
    }

    producer = rd_kafka_new(RD_KAFKA_PRODUCER, conf, errstr, sizeof(errstr));
    if (!producer)
    {
        g_error("Failed to create new producer: %s", errstr);
        logError("Failed to create new producer!", true);
        return 1;
    }

    conf = NULL;
    return 0;
}

注:消息回调函数dr_msg_cb仅用来报告发送消息时的可能错误。

消息发送代码:

int sendKafkaMessage(char *kafkaMessage)
{
    int message_count = 1;
    const char *topic = kafkaTopic;
    const char *value = kafkaMessage;

    for (int i = 0; i < message_count; i++)
    {
        size_t value_len = strlen(value);
        rd_kafka_resp_err_t err;

        err = rd_kafka_producev(producer,
            RD_KAFKA_V_TOPIC(topic),
            RD_KAFKA_V_MSGFLAGS(RD_KAFKA_MSG_F_COPY),
            RD_KAFKA_V_KEY(NULL, 0),
            RD_KAFKA_V_VALUE((void*)value, value_len),
            RD_KAFKA_V_OPAQUE(NULL),
            RD_KAFKA_V_END);

        if (err != RD_KAFKA_RESP_ERR_NO_ERROR)
        {
            // g_warning("Failed to produce to topic %s: %s", topic, rd_kafka_err2str(err));
            // logError("Failed to produce topic!", true);
            return 1;
        }
        else
        {
            // g_message("Produced event to topic %s: value = %12s", topic, value);
        }

        rd_kafka_poll(producer, 0);
    }

    // g_message("Flushing final messages..");
    rd_kafka_flush(producer, 100);

    if (rd_kafka_outq_len(producer) > 0)
    {
        // g_warning("%d message(s) were not delivered", rd_kafka_outq_len(producer));
        // logError("Kafka message(s) were not delivered!", true);
        return 1;
    }

    // g_message("%d events were produced to topic %s.", message_count, topic);
    return 0;
}

Valgrind检测结果:

==19032== 92,178 bytes in 9 blocks are definitely lost in loss record 45 of 45
==19032==    at 0x4848899: malloc (in /usr/libexec/valgrind/vgpreload_memcheck-amd64-linux.so)
==19032==    by 0x4A37F15: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x49FC06A: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x49E48E3: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x49F0B59: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x49F0F79: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x49B0D67: ??? (in /home/.../build/libs/librdkafka.so.1)
==19032==    by 0x4DC7934: start_thread (pthread_create.c:439)
==19032==    by 0x4E58BF3: clone (clone.S:100)

从Valgrind的输出看,泄漏点似乎在librdkafka库内部,但我不确定是不是自己的用法有问题导致的——比如有没有正确释放资源?或者在初始化、发送消息的流程中有没有遗漏什么步骤?希望有经验的朋友能帮我分析一下,谢谢!

备注:内容来源于stack exchange,提问作者Sebastian

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.21 08:13:08