- public EndlessConsumer<String, Long> endlessConsumer(
- org.apache.kafka.clients.consumer.Consumer<String, Long> kafkaConsumer,
- ExecutorService executor,
- Consumer<ConsumerRecord<String, Long>> handler,
- KafkaProperties kafkaProperties,
- ApplicationProperties applicationProperties)
+ public ProducerFactory<String, Object> producerFactory(KafkaProperties properties) {
+ return new DefaultKafkaProducerFactory<>(
+ properties.getProducer().buildProperties(),
+ new StringSerializer(),
+ new DelegatingByTypeSerializer(Map.of(
+ byte[].class, new ByteArraySerializer(),
+ ClientMessage.class, new JsonSerializer<>())));
+ }
+
+ @Bean
+ public KafkaTemplate<String, Object> kafkaTemplate(
+ ProducerFactory<String, Object> producerFactory) {
+
+ return new KafkaTemplate<>(producerFactory);
+ }
+
+ @Bean
+ public DeadLetterPublishingRecoverer recoverer(
+ ApplicationProperties properties,
+ KafkaOperations<?, ?> template)