rd_kafka_offsets_for_times()返回-186错误的代码修正方法
报错原因
你遇到的-186: Local: Invalid argument or configuration错误由两个错误用法直接导致:
- 手动操作分区列表的内部成员:
rd_kafka_topic_partition_list_t的elems、cnt都是结构内部维护的字段,你手动创建partition对象后直接赋值给elems,不会更新列表内部的元素计数cnt(此时cnt还是0),接口校验时发现传入的列表无有效元素,直接返回参数错误。 - 接口用途理解偏差:
rd_kafka_offsets_for_times()是按时间戳查询对应偏移量的接口,如果你要获取高水位(最新)偏移量,需要将对应分区的offset字段设为特殊值RD_KAFKA_OFFSET_END,填0代表查询1970年1月1日0点时间戳对应的偏移量,和你要获取高水位的预期不符。
修正方案
- 不要手动创建
rd_kafka_topic_partition_t对象、不要直接修改列表的内部指针,通过官方提供的rd_kafka_topic_partition_list_add接口往列表里添加分区,接口会自动维护列表的内存和计数。 - 根据你的查询目标设置正确的offset值:查高水位填
RD_KAFKA_OFFSET_END,查最早偏移填RD_KAFKA_OFFSET_BEGINNING,查指定时间点的偏移填对应毫秒级Unix时间戳。 - 接口调用成功后,除了检查接口整体返回值,还要逐个检查列表内每个分区的
err字段,避免单分区查询失败未被捕获。 - 用完分区列表后记得调用销毁接口释放内存,避免泄漏。
修正后的代码
// 创建初始容量为1的分区列表 rd_kafka_topic_partition_list_t* partition_list = rd_kafka_topic_partition_list_new(1); // 向列表添加目标主题的0号分区,接口返回对应分区对象的指针 rd_kafka_topic_partition_t* pt0 = rd_kafka_topic_partition_list_add(partition_list, topic, 0); // 设置查询目标为分区高水位(最新偏移量) pt0->offset = RD_KAFKA_OFFSET_END; rd_kafka_resp_err_t offsets_err = rd_kafka_offsets_for_times(rk, partition_list, 5000); if (offsets_err != RD_KAFKA_RESP_ERR_NO_ERROR) { printf("ERROR: Failed to get offsets: %d: %s.\n", offsets_err, rd_kafka_err2str(offsets_err)); } else { // 遍历读取每个分区的查询结果 for (int i = 0; i < partition_list->cnt; i++) { rd_kafka_topic_partition_t* p = &partition_list->elems[i]; if (p->err != RD_KAFKA_RESP_ERR_NO_ERROR) { printf("Partition %d query failed: %s\n", p->partition, rd_kafka_err2str(p->err)); continue; } printf("Partition %d target offset: %lld\n", p->partition, p->offset); } } // 销毁分区列表,释放内存 rd_kafka_topic_partition_list_destroy(partition_list);
补充提示
如果你只需要查询单分区的高低水位偏移量,也可以使用专门的rd_kafka_query_watermark_offsets()接口,不需要构造分区列表,调用逻辑更简单。所有librdkafka提供的集合类结构,都要使用配套的增删改接口操作,不要直接修改内部成员,否则很容易触发参数错误或者内存崩溃。
内容的提问来源于stack exchange,提问作者hermit.crab
相关产品推荐
相关产品推荐

