public class Demo001Producer {
public static void main(String[] args) {
//创建构建参数
Map<String,Object> map = new HashMap<>();
map.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");
map.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());//序列化类型
map.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());//序列化类型
//创建生产者对象
KafkaProducer<String,String> producer = new KafkaProducer<>(map);
//创建数据
ProducerRecord<String,String> record = new ProducerRecord<>("test", "wyl", "hello kafka");
//发送数据到kafka-topic
producer.send(record);
//关闭生产者对象
producer.close();
}
}
public class Demo001Consumer {
public static void main(String[] args) {
//配置map
Map<String, Object> configs = new HashMap<>();
configs.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
configs.put(ConsumerConfig.GROUP_ID_CONFIG, "test"); // 消费者组ID
configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 反序列化类型,字符串类型
configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 反序列化类型,字符串类型
//创建消费者
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(configs);
consumer.subscribe(Collections.singletonList("test"));
while (true) {
consumer.poll(1000).forEach(record -> {
System.out.println(record.key() + ":" + record.value());
});
}
}
}