У меня есть приложение Springboot, которое подключается к брокеру Kafka и дает сообщения для настроенной темы Kafka. У меня есть, прежде всего, SenderComponent, который планирует задачи, основанные на выражении Cron.package com.sporting.scheduler.sender.application;
import com.sporting.scheduler.registry.domain.OutgoingRouting;
import com.sporting.scheduler.registry.domain.ScheduledJobRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.scheduling.TaskScheduler;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.scheduling.support.CronTrigger;
import org.springframework.stereotype.Component;
import java.util.HashSet;
import java.util.Set;
@Component
@RequiredArgsConstructor
@EnableScheduling
@Slf4j
public class SenderComponent {
static final String SENDER_PREFIX = "[SENDER] - ";
private final Set scheduledRoutings = new HashSet();
@Value("${sporting.kafka.url}")
private String url;
@Value("${sporting.kafka.sender-topic}")
private String topicName;
private final ScheduledJobRepository repository;
private final TaskScheduler scheduler;
private final SenderKafkaProducer senderKafkaProducer;
@Bean
public CommandLineRunner start() {
return args -> log.info("{}The scheduler sender started, listening for scheduled tasks to trigger.", SENDER_PREFIX);
}
@Scheduled(fixedRate = 1000 * 60, initialDelay = 1000) // Every minute, wait a second to boot
public void senderTriggers() {
Set routings = repository.fetchOutgoingRoutings();
if (routings.isEmpty()) {
log.info("{}No outgoing routings found in the job registry.", SENDER_PREFIX);
return;
}
routings.forEach(this::scheduleTask);
}
private void scheduleTask(OutgoingRouting outgoingRouting) {
if (!scheduledRoutings.contains(outgoingRouting.getCode())) {
scheduler.schedule(() -> senderKafkaProducer.send(topicName, outgoingRouting.getCode()), new CronTrigger(outgoingRouting.getCronExpression()));
scheduledRoutings.add(outgoingRouting.getCode());
}
}
}
< /code>
Это код, который я записал в senderkafkaproducer < /p>
package com.sporting.scheduler.sender.application;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.stereotype.Component;
import java.util.HashMap;
import java.util.Map;
@Component
@RequiredArgsConstructor
@Slf4j
public class SenderKafkaProducer {
@Bean
public ProducerFactory producerFactory() {
return new DefaultKafkaProducerFactory(producerConfigs());
}
@Bean
public Map producerConfigs() {
Map props = new HashMap();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_METRICS_PUSH_CONFIG, "false");
return props;
}
@Bean
public KafkaTemplate kafkaTemplate() {
return new KafkaTemplate(producerFactory());
}
public void send(String topic, String message) {
kafkaTemplate().send(topic, message);
log.info("Sent message [{}] to topic [{}]", message, topic);
}
}
< /code>
Я вижу, что я получаю следующие строки журнала < /p>
2025-08-07T16:52:00.015+02:00 INFO 32399 --- [Sender] [ scheduling-1] c.t.s.d.a.SenderKafkaProducer : Sent message [BINGO_MAN] to topic [SCH-list]
< /code>
Так что я предполагаю, что брокер Kafka хорошо получил мое сообщение. Тем временем я использую потребителя Kafka на этом брокере, но он вообще не получает никакого сообщения. Это дает мне следующий выход < /p>
Metadata for SCH-list (from broker 1: localhost:9092/1):
1 brokers:
broker 1 at localhost:9092 (controller)
1 topics:
topic "SCH-list" with 1 partitions:
partition 0, leader 1, replicas: 1, isrs: 1
< /code>
Но мое сообщение bingo_man не принимается. У меня есть идея, что мое приложение Springboot правильно подключается к брокеру Kafka, но что -то идет не так с продюсером.
Может ли кто -нибудь указать мне, что я делаю неправильно, а также дать советы, как я могу проверить в своем коде, если сообщение было хорошо получено моим брокером Kafka?
>
Подробнее здесь: https://stackoverflow.com/questions/797 ... kaproducer