使用avro生成entity文件可以查看这篇文章https://cloud.tencent.com/developer/article/1683778
生产者代码
代码语言:javascript复制 public static void CustomerTest() {
Properties kafkaProps = new Properties();
kafkaProps.put("bootstrap.servers","192.168.0.31:9092,192.168.0.32:9092,192.168.0.33:9092");
kafkaProps.put("key.serializer","org.apache.kafka.common.serialization.StringSerializer");
kafkaProps.put("value.serializer","org.apache.kafka.common.serialization.ByteArraySerializer");
KafkaProducer producer = new KafkaProducer<String,byte[]>(kafkaProps);
for(int i = 0;i < 1000;i ){
Customer customer = new Customer();
customer.setEmail("23132@163.com-" i);
customer.setName("ric-" i);
customer.setId(i);
customer.setImages(null);
ByteArrayOutputStream out = new ByteArrayOutputStream();
BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(out, (BinaryEncoder)null);
SpecificDatumWriter writer = new SpecificDatumWriter(customer.getSchema());
try {
writer.write(customer, encoder);
encoder.flush();
out.close();
} catch (IOException e) {
e.printStackTrace();
}
ProducerRecord<String,byte[]> record = new ProducerRecord<String, byte[]>("Customer","customer-" i,out.toByteArray());
producer.send(record);
}
producer.close();
}
消费者代码
代码语言:javascript复制 public static void CustomerTest() {
Properties kafkaProps = new Properties();
kafkaProps.put("bootstrap.servers","192.168.0.31:9092,192.168.0.32:9092,192.168.0.33:9092");
kafkaProps.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer");
kafkaProps.put("value.deserializer","org.apache.kafka.common.serialization.ByteArrayDeserializer");
kafkaProps.put("group.id","DemoAvroKafkaConsumer2");
kafkaProps.put("auto.offset.reset","earliest");
KafkaConsumer<String ,byte[]> consumer = new KafkaConsumer<String, byte[]>(kafkaProps);
consumer.subscribe(Collections.singletonList("Customer"));
SpecificDatumReader<Customer> reader = new SpecificDatumReader<>(Customer.getClassSchema());
try {
while (true){
ConsumerRecords<String,byte[]> records = consumer.poll(10);
for(ConsumerRecord<String,byte[]> record : records){
Decoder decoder = DecoderFactory.get().binaryDecoder(record.value(), null);
Customer customer = null;
try {
customer = reader.read(null,decoder);
System.out.println(record.key() ":" customer.get("id") "t" customer.get("name") "t" customer.get("email"));
} catch (IOException e) {
e.printStackTrace();
}
}
}
} finally {
consumer.close();
}
}
相关pom依赖
代码语言:javascript复制 <dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.11</artifactId>
<version>1.0.0</version>
</dependency>
<dependency>
<groupId>org.apache.avro</groupId>
<artifactId>avro</artifactId>
<version>1.8.2</version>
</dependency>
<dependency>
<groupId>org.apache.avro</groupId>
<artifactId>avro-tools</artifactId>
<version>1.8.2</version>
</dependency>
<dependency>
<groupId>com.twitter</groupId>
<artifactId>bijection-avro_2.11</artifactId>
<version>0.9.6</version>
</dependency>