springboot整合kafka发送同步消息
技术标签: kafka kafka
1.引入kafka的maven依赖
<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>2.1.0.RELEASE</version> </dependency>
2.创建topic
./bin/kafka-topics.sh --create --zookeeper 192.168.211.137:2181 --partitions 3 --replication-factor 1 \ --topic syn-message
3.代码
package com.wyu.tt06kafkademo.demo; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import java.util.Properties; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; /** * 同步发送消息 */ public class SynProducer { private static Properties getProps(){ Properties props = new Properties(); props.put("bootstrap.servers", "192.168.211.137:9092"); props.put("acks", "all"); // 发送所有ISR props.put("retries", 2); // 重试次数 props.put("batch.size", 16384); // 批量发送大小 props.put("buffer.memory", 102400); // 缓存大小,根据本机内存大小配置 props.put("linger.ms", 1000); // 发送频率,满足任务一个条件发送 props.put("client.id", "producer-syn-1"); // 发送端id,便于统计 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); return props; } public static void main(String[] args) { KafkaProducer<String, String> producer = new KafkaProducer<>(getProps()); for(int i=0; i< 10; i++){ // 三个参数,topic,key:用户分配partition,value:发送的值 ProducerRecord<String, String> record = new ProducerRecord<>("syn-message", "topic_"+i,"message-"+i); Future<RecordMetadata> metadataFuture = producer.send(record); RecordMetadata recordMetadata = null; try { recordMetadata = metadataFuture.get(); System.out.println("发送成功!"); System.out.println("topic:"+recordMetadata.topic()); System.out.println("partition:"+recordMetadata.partition()); System.out.println("offset:"+recordMetadata.offset()); } catch (InterruptedException|ExecutionException e) { System.out.println("发送失败!"); e.printStackTrace(); } } producer.flush(); producer.close(); } }
4.消费syn-message的数据
bin/kafka-console-consumer.sh --bootstrap-server 192.168.211.137:9092 --from-beginning --topic syn-message