【问题标题】:Consuming multiple Kafka topics in a Rust consumer在 Rust 消费者中使用多个 Kafka 主题
【发布时间】:2015-10-28 04:23:08
【问题描述】:

我正在尝试使用 Rust 消费者来读取多个主题。这是我现在的代码:

extern crate kafka;
use kafka::client::KafkaClient;
use kafka::consumer::Consumer;
use kafka::utils;
fn main(){
    let mut client = KafkaClient::new(vec!("localhost:9092".to_owned()));
    let res = client.load_metadata_all();
    let topics = client.topic_partitions.keys().cloned().collect(); 
    let offsets = client.fetch_offsets(topics, -1);
    for topic in &topics {
    let mut con = Consumer::new(client, "test-consumer-group".to_owned(), "topic".to_owned()).partition(0);
    let mut messagedata = 0;
    for msg in con {
        println!("{}", str::from_utf8(&msg.message).unwrap().to_string());
    }
  }
}

以下是错误:

    src/main.rs:201:19: 201:25 error: use of moved value: `topics` [E0382]
src/main.rs:201     for topic in &topics {
                                  ^~~~~~
    note: in expansion of for loop expansion
    src/main.rs:201:5: 210:6 note: expansion site
    src/main.rs:167:40: 167:46 note: `topics` moved here because it has type `collections::vec::Vec<collections::string::String>`, which is non-copyable
    src/main.rs:167     let offsets = client.fetch_offsets(topics, -1);
                                                       ^~~~~~
    src/main.rs:203:37: 203:43 error: use of moved value: `client` [E0382]
    src/main.rs:203     let mut con = Consumer::new(client, "test-consumer-group".to_owned(), "topicname".to_owned()).partition(0);
                                                    ^~~~~~
    note: in expansion of for loop expansion
    src/main.rs:201:5: 210:6 note: expansion site
    note: `client` was previously moved here because it has type     `kafka::client::KafkaClient`, which is non-copyable
    error: aborting due to 2 previous errors

为了更好地解释我的问题,这是我仅针对一个主题的部分可行代码:

let mut con = Consumer::new(client, "test-consumer-group".to_owned(), "testtopic".to_owned()).partition(0);

for msg in con {
    println!("{}", str::from_utf8(&msg.message).unwrap().to_string());
}

我测试了 fetch_message 函数,它适用于多个主题,但我得到的结果(msgs)是 Topicmessage,我不知道如何从 Topicmessage 获取消息。

let msgs = client.fetch_messages_multi(vec!(utils::TopicPartitionOffset{
                                            topic: "topic1".to_string(),
                                            partition: 0,
                                            offset: 0 //from the begining
                                            },
                                        utils::TopicPartitionOffset{
                                            topic: "topic2".to_string(),
                                            partition: 0,
                                            offset: 0
                                        },
                                        utils::TopicPartitionOffset{
                                            topic: "topic3".to_string(),
                                            partition: 0,
                                            offset: 0
                                        }));
for msg in msgs{
    println!("{}", msg);
}

【问题讨论】:

  • @Shepmaster 感谢您的提醒,我修改了我的问题。
  • 这些问题实际上是 Rust 的基础:所有权。请阅读其他有 use of moved value 错误的 stackoverflow 问题,并阅读本书关于所有权的章节:doc.rust-lang.org/stable/book/ownership.html
  • @ker 谢谢你的帮助!

标签: rust apache-kafka kafka-consumer-api


【解决方案1】:

最后,我像这样更改了代码,它可以工作了。

let msgs = client.fetch_messages_multi(...).unwrap();
for msg in msgs{
     println!("{}", msg.message);
}

【讨论】:

  • 请解释问题中的问题以及为什么您的解决方案是正确的
猜你喜欢
  • 2017-01-26
  • 1970-01-01
  • 1970-01-01
  • 2018-12-31
  • 2018-08-31
  • 1970-01-01
  • 2020-09-19
  • 1970-01-01
  • 2020-01-05
相关资源
最近更新 更多