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

