【发布时间】: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