Kafka整合WebFlux實(shí)踐
更新時(shí)間:2026年03月06日 15:25:28 作者:西紅柿系番茄
文章介紹了如何在Kafka中整合WebFlux,包括引入依賴和代碼示例,并對(duì)相關(guān)知識(shí)進(jìn)行了總結(jié)
Kafka整合WebFlux
1、引入依賴
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor.kafka</groupId>
<artifactId>reactor-kafka</artifactId>
<version>1.1.0.RELEASE</version>
</dependency>2、代碼示例
@Component
public class KafkaService {
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
private KafkaSender<String, String> kafkaSender;
private KafkaReceiver<String, String> kafkaReceiver;
@PostConstruct
public void init() {
final Map<String, Object> producerProps = new HashMap<>();
producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
final SenderOptions<String, String> producerOptions = SenderOptions.create(producerProps);
this.kafkaSender = KafkaSender.create(producerOptions);
final Map<String, Object> consumerProps = new HashMap<>();
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG, "payment-validator-1");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "payment-validator");
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
ReceiverOptions<String, String> consumerOptions = ReceiverOptions.<String, String>create(consumerProps)
.subscription(Collections.singleton("demo"))
.addAssignListener(partitions -> System.out.println("onPartitionsAssigned " + partitions))
.addRevokeListener(partitions -> System.out.println("onPartitionsRevoked " + partitions));
KafkaReceiver<String, String> kafkaReceiver = KafkaReceiver.create(consumerOptions);
kafkaReceiver.receive().doOnNext(r -> {
System.out.println(r.value());
r.receiverOffset().acknowledge();
}).subscribe();
this.kafkaReceiver = kafkaReceiver;
}
public Mono< ?> send() {
SenderRecord<String, String, Object> senderRecord = SenderRecord.create(new ProducerRecord<>("demo", value()), 1);
return kafkaSender.send(Mono.just(senderRecord)).next();
}
private String value() {
Map<String, String> map = new HashMap<>();
map.put("name", UUID.randomUUID().toString());
try {
return OBJECT_MAPPER.writeValueAsString(map);
} catch (JsonProcessingException e) {
return "{}";
}
}
}3、其它
server:
port: 8888
spring:
jackson:
serialization:
FAIL_ON_EMPTY_BEANS: false總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
eclipse報(bào)錯(cuò) eclipse啟動(dòng)報(bào)錯(cuò)解決方法
本文將介紹eclipse啟動(dòng)報(bào)錯(cuò)解決方法,需要了解的朋友可以參考下2012-11-11
java8 實(shí)現(xiàn)提取集合對(duì)象的每個(gè)屬性
這篇文章主要介紹了java8 實(shí)現(xiàn)提取集合對(duì)象的每個(gè)屬性方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過來看看吧2021-02-02
Java中的StopWatch計(jì)時(shí)利器使用指南
StopWatch通常用于測(cè)量一段代碼執(zhí)行所花費(fèi)的時(shí)間,它能夠精確地記錄開始時(shí)間、結(jié)束時(shí)間,并計(jì)算出這中間的時(shí)間差,下面給大家介紹Java中的StopWatch計(jì)時(shí)利器的深度解析與使用指南,感興趣的朋友一起看看吧2025-05-05
Mybatis 動(dòng)態(tài)SQL的幾種實(shí)現(xiàn)方法
這篇文章主要介紹了Mybatis 動(dòng)態(tài)SQL的幾種實(shí)現(xiàn)方法,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-11-11
Java Kafka實(shí)現(xiàn)延遲隊(duì)列的示例代碼
kafka作為一個(gè)使用廣泛的消息隊(duì)列,很多人都不會(huì)陌生。本文將利用Kafka實(shí)現(xiàn)延遲隊(duì)列,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以嘗試一下2022-08-08

