使用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
相关产品推荐
相关产品推荐

