【发布时间】:2018-11-08 16:09:34
【问题描述】:
我已经编写了这个 Kafka 生产者并从桌面读取一个文件,然后将文件中的数据作为值推送,并通过每次读取每一行时添加一个来生成我自己的密钥。这是正确的方法还是我做了我不应该做的事情?请需要一些建议。 我可以在我的主题中看到消息,但每个消息都与一个键相关,所以如果我有一个用例,如果我从外部读取它,我可以推送任何这样的日志数据。我可以使用日志数据作为值,还是应该采用完全不同的逻辑。 请帮忙
import java.io.BufferedReader;
import java.io.File;
import java.io.FileReader;
import java.io.IOException;
import java.util.Properties;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
public class SyncProducer {
public static void main(String[] args) throws IOException {
File file = new File("/Users/adityaverma/Desktop/ParseData.txt");
BufferedReader br = new BufferedReader(new FileReader(file));
Properties properties = new Properties();
properties.setProperty("bootstrap.servers","127.0.0.1:9092");
properties.setProperty("key.serializer",StringSerializer.class.getName()); // our key and values are String
properties.setProperty("value.serializer",StringSerializer.class.getName());
properties.setProperty("acks", "1");
properties.setProperty("retries", "3");
properties.setProperty("linger.ms", "1");
Producer<String,String> producer = new org.apache.kafka.clients.producer.KafkaProducer<String,String>(properties);
// these will go in random partition as we increment the key
String line = " ";
int key = 0;
while((line = br.readLine()) != null){
// System.out.println(line);
ProducerRecord<String,String> producerRecord = new ProducerRecord<String,String>("try_Buffered3Part",Integer.toString(key),line);
key++;
System.out.println(key);
producer.send(producerRecord);
}
producer.close();
System.out.println("exit");
}
}
【问题讨论】:
标签: java file apache-kafka hadoop-yarn hadoop2