【问题标题】:How to save only last message from topic to a file?如何仅将主题中的最后一条消息保存到文件中?
【发布时间】:2018-06-17 03:31:57
【问题描述】:

如标题 - 如何仅将主题中的最后一条消息保存到文件中?我试过这样做,但是当我打开一个文件时,它包含我之前发送的大约最后 10 条消息。这是我的代码:

public class Consumer {
FileWriter fw;
PrintWriter pw;

public void Consume(String topic){
    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("group.id", "ranking_consumer");
    props.put("enable.auto.commit", "false");
    props.put("auto.offset.reset", "latest");
    props.put("auto.commit.interval.ms", "1000");
    props.put("session.timeout.ms", "30000");
    props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
    consumer.subscribe(Arrays.asList(topic));

    File dir = new File("C:/kafka-logs/kafka-logs-0/data");
    if (!dir.exists()) {
        dir.mkdir();
    }

    try{
    while(true){
        ConsumerRecords<String, String> records = consumer.poll(1000);
        int chunkSize = records.count();
        int recordIndex = 0;

        for(ConsumerRecord<String, String> record : records){
            recordIndex++;
            if (recordIndex == chunkSize){
            fw = new FileWriter("C:/kafka-logs/kafka-logs-0/data/msg-"+topic, false);
            pw = new PrintWriter(fw);
            pw.print(record.value());
            pw.close();
            }
        }
    }} catch(IOException e){
        e.printStackTrace();
    }finally{
    consumer.close();
    }
}

}

注意:添加了 dbustosp 提出的代码行

注意 2: 更清楚地解释我的问题。我有一个用于生产者的简单 JFrame,带有“选择文件”按钮,它将 .csv 文件作为字符串变量发送。然后我有第二个简单的 JFrame 用于接收最后发布的消息并将其保存到文件中。也许我有一个错误的消费者线程声明,但是当我从生产者发送几条消息时,我收到一个包含一些相同消息的文件,而不是只有一个 - 最后一个。下面我将我的代码发送给消费者的 JFrame 和生产者:

public class Producer {

public void Produce(String msg, String topic){

    Properties props = new Properties();
    props.put("bootstrap.servers", "localhost:9092");
    props.put("acks", "all");
    props.put("retries", 0);
    props.put("batch.size", 16384);
    props.put("linger.ms", 1);
    props.put("buffer.memory", 33554432);
    props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

    KafkaProducer<String, String> producer = new KafkaProducer<>(props);
    ProducerRecord<String, String> record = new ProducerRecord<>(topic, null, msg);
    producer.send(record);
    producer.close();
}
}

public class Kafka_Consumer extends javax.swing.JFrame {

Thread t;

public Kafka_Consumer() {
    initComponents();
}
private void btn_ReceiveFromKafkaActionPerformed(java.awt.event.ActionEvent evt) {                                                     
    if (textField_topic.getText().equals("")) {
        showMessageDialog(null, "Nie wprowadziłeś/aś tematu!");
    }
    else{
        Consumer consumer = new Consumer();
        t = new Thread(new Runnable(){
        public void run(){
        consumer.Consume(textField_topic.getText());
        }
    });
    t.start();
    }

}

【问题讨论】:

  • 您是否注意到每次通过“else”子句时都在创建一个新的消费者。此外,还初始化并执行了一个新线程。您根本没有停止线程。您遇到的问题是您如何处理消费者和线程。您可以围绕这个问题提出一个新问题。

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


【解决方案1】:

我没有看到任何代码在轮询请求后处理最后一条消息。 在您的代码中进行以下修改应该可以工作。

// Getting the chunk size got from poll
int chunkSize = records.count();
// Index to count the number of the message
int recordIndex = 0;
for(ConsumerRecord<String, String> record : records) {
    // Increasing to count the number of message visited
    recordIndex++;
    // Checking the last message
    if (recordIndex == size) {
        fw = new FileWriter("C:/kafka-logs/kafka-logs-0/data/msg-"+topic, false);
        pw = new PrintWriter(fw);
        pw.print(record.value());
        pw.close();
    }
}

【讨论】:

  • 第一次成功了,但是当我添加了几条消息时,仍然有一些相同的消息
  • 那段代码正在选择您从轮询请求中带来的块中的最后一条记录,我看不出有任何理由不做您最初想要的工作。如果您可以添加更多详细信息,说明您希望看到什么,队列中已经有什么数据,以及是否有任何进程将数据推送到您正在使用的主题/分区中,那就太好了。
  • 添加注释2以获得更好的解释
猜你喜欢
  • 1970-01-01
  • 2018-04-20
  • 2011-01-14
  • 1970-01-01
  • 2020-11-21
  • 1970-01-01
  • 2019-11-14
  • 2022-01-17
  • 2018-07-11
相关资源
最近更新 更多