import org.springframework.util.concurrent.ListenableFuture;
-
@Slf4j
@SpringBootApplication
public class Application implements ApplicationRunner
@Autowired
KafkaTemplate<String, String> kafkaTemplate;
-
- void send(String key, String value)
- {
- ListenableFuture<SendResult<String, String>> listenableFuture =
- kafkaTemplate.send("test", key, value);
-
- listenableFuture.addCallback(
- result -> log.info(
- "Sent {}={} to partition={}, offset={}",
- result.getProducerRecord().key(),
- result.getProducerRecord().value(),
- result.getRecordMetadata().partition(),
- result.getRecordMetadata().offset()),
- e -> log.error("ERROR sendig message", e));
- }
-
@Override
public void run(ApplicationArguments args)
{
for (int i = 0; i < 100; i++)
{
- send(Long.toString(i%10), Long.toString(i));
+ ListenableFuture<SendResult<String, String>> listenableFuture =
+ kafkaTemplate.send("test", Long.toString(i%10), Long.toString(i));
+
+ listenableFuture.addCallback(
+ result -> log.info(
+ "Sent {}={} to partition={}, offset={}",
+ result.getProducerRecord().key(),
+ result.getProducerRecord().value(),
+ result.getRecordMetadata().partition(),
+ result.getRecordMetadata().offset()),
+ e -> log.error("ERROR sendig message", e));
}
}
-
public static void main(String[] args)
{
SpringApplication.run(Application.class, args);