使用modern-cpp-kafka与librdkafka销毁Headers时崩溃问题排查
问题:Kafka Headers内存分配导致
rd_kafka_headers_destroy崩溃 场景与代码
我在本地函数中创建kafka::Headers,其值通过堆动态分配内存生成,但未手动释放指针:
long n1; long n2; kafka::Headers headers; auto pN1Bytes = (decltype(n1)*)malloc(sizeof(n1)); *pN1Bytes = n1; auto pN2Bytes = (decltype(n2)*)malloc(sizeof(n2)); *pN2Bytes = n2; headers.push_back(kafka::Header(kafka::Header::Key("n1"), kafka::Header::Value(pN1Bytes, sizeof(n1)))); headers.push_back(kafka::Header(kafka::Header::Key("n2"), kafka::Header::Value(pN2Bytes, sizeof(n2))));
消息发送通过以下方法实现:
bool KafkaProducerService::Produce(const kafka::Key& key, const std::string& msg, const std::string& topic, const kafka::Headers& headers) { auto ret = true; try { auto record = kafka::clients::producer::ProducerRecord(topic, mPartition, key, kafka::Value(msg.c_str(), msg.size())); record.headers() = headers; mpProducer->send( record, mProduceCallback, kafka::clients::producer::KafkaProducer::SendOption::ToCopyRecordValue, kafka::clients::producer::KafkaProducer::ActionWhileQueueIsFull::NoBlock ); mpProducer->flush(std::chrono::milliseconds(50)); } catch (const kafka::KafkaException& e) { std::cerr << "% Unexpected exception caught: " << e.what() << std::endl; ret = false; } return ret; }
崩溃情况
程序运行时删除Kafka服务器上的对应Topic后,在rd_kafka_headers_destroy函数处崩溃。
版本信息
- modern-cpp-kafka: v2024.07.03
- librdkafka: 2.4.0
疑问
原本认为librdkafka会自动释放Headers,但程序却崩溃了,请问原因是什么?
解答
崩溃的核心原因是Headers的内存所有权冲突与双重释放,具体拆解:
modern-cpp-kafka的Headers内存管理逻辑
kafka::Headers是librdkafka中rd_kafka_headers_t的RAII封装,当kafka::Headers对象被销毁时,会自动调用rd_kafka_headers_destroy释放底层资源。但你创建kafka::Header::Value时传入的是手动malloc的堆指针,modern-cpp-kafka默认不会复制这个指针指向的内存,也不会接管其所有权——它只是将指针直接传递给底层的librdkafka结构。所有权冲突与双重释放的触发
- 你在本地函数中创建的
headers变量,在调用Produce方法后,随着本地函数退出会被销毁,此时会触发rd_kafka_headers_destroy释放一次Headers资源。 - 虽然你调用了
mpProducer->flush(50ms),但这个超时时间很短,当Topic被删除后,消息发送流程可能还在后台处理。后续librdkafka在清理发送失败的消息资源时,会再次尝试调用rd_kafka_headers_destroy处理同一个已经被释放的rd_kafka_headers_t指针,导致双重释放,直接触发崩溃。
ToCopyRecordValue的局限性
你使用的ToCopyRecordValue选项仅针对消息体(Record Value)的内存拷贝,不会对Headers的内存进行复制。因此Headers的底层指针仍然会被librdkafka持有,当原kafka::Headers对象销毁后,就会出现野指针或重复释放的问题。
修复建议
- 让modern-cpp-kafka接管内存所有权:使用
kafka::Header::Value的带释放回调的重载,或者直接使用值拷贝的方式(比如将栈上的n1、n2直接传入,让modern-cpp-kafka自动复制内存,而不是手动分配堆内存)。
示例:kafka::Headers headers; // 方式1:指定不自动释放,后续手动管理(不推荐) headers.push_back(kafka::Header("n1", kafka::Value(&n1, sizeof(n1), kafka::Value::Deallocator::None))); // 方式2:让库自动复制内容(推荐) headers.push_back(kafka::Header("n1", std::string(reinterpret_cast<char*>(&n1), sizeof(n1)))); - 延长
kafka::Headers的生命周期:确保headers对象在消息发送完成(包括回调处理)前不会被销毁,避免提前释放底层资源。 - 避免手动分配内存:尽量使用栈内存或让modern-cpp-kafka管理内存,减少手动内存分配带来的所有权问题。
内容的提问来源于stack exchange,提问作者TH Hu
相关产品推荐
相关产品推荐

