在Kafka上异步发送数据
對于一個項目,我試圖記錄用戶的基本交易,例如添加和刪除一個項目以及多種類型的項目,并為每筆交易向kafka發送一條消息。 日志機制的準確性不是至關重要的,在kafka服務器停機的情況下,我不希望它阻止我的業務代碼。 在這種情況下,將數據發送到kafka的異步方法是一種更好的方法。
我的kafka生產者代碼在其引導項目中。 為了使其異步,我只需要添加兩個注釋:@EnableAsync和@Async。
@EnableAsync將在您的配置類中使用(還要記住,帶有@SpringBootApplication的類也是配置類),并將嘗試查找TaskExecutor bean。 如果沒有,它將創建一個SimpleAsyncTaskExecutor。 SimpleAsyncTaskExecutor適用于玩具項目,但對于任何大于此的項目都存在一定的風險,因為它不限制并發線程,也不會重用線程。 為了安全起見,我們還將添加一個任務執行者bean。
所以,
@SpringBootApplication public class KafkaUtilsApplication { public static void main(String[] args) { SpringApplication.run(KafkaUtilsApplication. class , args); } }會變成
@EnableAsync @SpringBootApplication public class KafkaUtilsApplication { public static void main(String[] args) { SpringApplication.run(KafkaUtilsApplication. class , args); } @Bean public Executor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize( 2 ); executor.setMaxPoolSize( 2 ); executor.setQueueCapacity( 500 ); executor.setThreadNamePrefix( "KafkaMsgExecutor-" ); executor.initialize(); return executor; } }如您所見,這里沒有太多變化。 我設置的默認值應根據您的應用程序需求進行調整。
我們需要的第二件事是添加@Async。
我的舊代碼是:
@Service public class KafkaProducerServiceImpl implements KafkaProducerService { private static final String TOPIC = "logs" ; @Autowired private KafkaTemplate<String, KafkaInfo> kafkaTemplate; @Override public void sendMessage(String id, KafkaType kafkaType, KafkaStatus kafkaStatus) { kafkaTemplate.send(TOPIC, new KafkaInfo(id, kafkaType, kafkaStatus); } }如您所見,同步代碼非常簡單。 它只需要kafkaTemplate并將消息對象發送到“ logs”主題。 我的新代碼比這更長。
@Service public class KafkaProducerServiceImpl implements KafkaProducerService { private static final String TOPIC = "logs" ; @Autowired private KafkaTemplate kafkaTemplate; @Async @Override public void sendMessage(String id, KafkaType kafkaType, KafkaStatus kafkaStatus) { ListenableFuture<SendResult<String, KafkaInfo>> future = kafkaTemplate.send(TOPIC, new KafkaInfo(id, kafkaType, kafkaStatus)); future.addCallback( new ListenableFutureCallback<>() { @Override public void onSuccess( final SendResult<String, KafkaInfo> message) { // left empty intentionally } @Override public void onFailure( final Throwable throwable) { // left empty intentionally } }); } }在這里,onSuccess()對我而言并不真正有意義。 但是onFailure()可以記錄異常,因此可以通知我我的kafka服務器是否存在問題。
我還要與您分享另一件事。 為了通過kafkatemplate發送對象,我必須為其配備序列化文件。
public class KafkaInfoSerializer implements Serializer<kafkainfo> { @Override public void configure(Map map, boolean b) { } @Override public byte [] serialize(String arg0, KafkaInfo info) { byte [] retVal = null ; ObjectMapper objectMapper = new ObjectMapper(); try { retVal = objectMapper.writeValueAsString(info).getBytes(); } catch (Exception e) { // log the exception } return retVal; } @Override public void close() { } }另外,不要忘記為其添加配置。 有幾種定義kafka的序列化器的方法。 最簡單的方法之一是將其添加到application.properties。
spring.kafka.producer.key-serializer = org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer = com.sezinkarli.kafkautils.serializer.KafkaInfoSerializer
現在,您有了一個啟動項目,該項目可以將異步對象發送到所需的主題。
翻譯自: https://www.javacodegeeks.com/2020/01/send-your-data-async-on-kafka.html
總結
以上是生活随笔為你收集整理的在Kafka上异步发送数据的全部內容,希望文章能夠幫你解決所遇到的問題。
 
                            
                        - 上一篇: H5画布
- 下一篇: 电脑自带游戏空当接龙(电脑小游戏空当接龙
