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

关于Kafka Consumer的subscription()函数用途及返回值异常的咨询

关于rust-rdkafka中Consumer::subscription()函数的疑问

背景说明

我使用Rust结合rust-rdkafka编写测试代码时,尝试调用Kafka Consumer的subscription()函数,但不清楚该函数的设计用途。查阅Kafka相关文档后未获得足够有效信息,且该函数底层基于C++的rdkafka实现,Java与Rust版本的文档内容一致。

调用代码及返回结果

调用代码:

let subscription = consumer.subscription().unwrap();

返回的数据如下:

element.topic: "kafka_test"
element.partition: -1
element.offset: Invalid
element.metadata: ""
element.topic: "kafka_test_1.json"
element.partition: -1
element.offset: Invalid
element.metadata: ""

核心困惑

我无法理解为何返回结果中的partition、offset处于无效状态,且metadata为空。我已多次调用poll()获取到有效记录,也尝试提交当前offset,但subscription()的返回结果始终没有变化。

测试代码片段

fn main() {
    // 每次程序运行时该值都会变化
    let group_id = "test_group_id" + run_counter;

    let mut client_config = ClientConfig::new();
    let client_config = client_config
        .set("bootstrap.servers", ...)
        .set("security.protocol", ...)
        .set("sasl.mechanisms", ...)
        .set("sasl.username", ...)
        .set("sasl.password", ...)
        .set("group.id", group_id)
        .set("enable.auto.commit", "false")
        .set("auto.offset.reset", "earliest")
        .set_log_level(rdkafka::config::RDKafkaLogLevel::Debug);

    let consumer: BaseConsumer = client_config
        .create()
        .expect("invalid producer config");
    
    consumer.subscribe(&topic_names).expect("consumer subscribe failed");

    test_function(&consumer, std::time::Duration::from_secs(30));
}

pub fn test_function(
    consumer: &BaseConsumer,
    poll_duration: Option<Duration>
) {
    let something = consumer.poll(poll_duration);
    log::info!("something={:?}", something); // 能获取到有效记录(不是"None")

    // 测试subscription()
    let subscription = consumer.subscription().unwrap();

    // 打印topic map的元素
    // 如前文所示,offset无效,partition无效……原因何在?
    let elements = subscription.elements();
    for element in elements.into_iter() {
        let topic = element.topic();
        let partition = element.partition();
        let offset = element.offset();
        let metadata = element.metadata();

        log::info!("element.topic: {:?}", topic);
        log::info!("element.partition: {:?}", partition);
        log::info!("element.offset: {:?}", offset);
        log::info!("element.metadata: {:?}", metadata);
    }
}

内容的提问来源于stack exchange,提问作者user2138149

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 17:05:40