логическая идея
Я настраивал приложение Springboot, которое создает сообщения для брокера Kafka, который работает в контейнере Docker.
Docker Setup of Kafka Broker
Я следовал за следующим шагом к установке Cafka Broker с Aronge Theme и Presser -Broker с Aromer с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry с Antry -Broker с Arner -Broker с Arting Arner -Broker. />docker run -d -p 9092:9020 --name broker apache/kafka:latest
docker exec --workdir /opt/kafka/bin/ -it broker sh
./kafka-topics.sh --bootstrap-server localhost:9092 --create --topic SCH-dispatcher
./kafka-console-consumer.sh --topic SCH-dispatcher --from-beginning --bootstrap-server localhost:9092
< /code>
code < /h1>
У меня есть следующие свойства, настроенные в моем приложении.# KAFKA CONFIGURATION
trax.kafka.address=localhost:9092
trax.kafka.dispatch-topic=SCH-dispatcher
< /code>
Рядом с этим у меня есть класс кафкаконфигурации, который создает компонент для создания соединения кафки. < /p>
package com.emea.toekeloo.shared.infrastructure;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaConfiguration {
@Value(value = "${trax.kafka.address}")
private String kafkaAddress;
@Bean
public ProducerFactory producerFactory() {
Map configProps = new HashMap();
configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaAddress);
configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.ENABLE_METRICS_PUSH_CONFIG, false);
return new DefaultKafkaProducerFactory(configProps);
}
@Bean
public KafkaTemplate kafkaTemplate() {
return new KafkaTemplate(producerFactory());
}
}
< /code>
Рядом с этим у меня есть класс производителя кафки: < /p>
package com.emea.toekeloo.dispatcher.application;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import java.util.concurrent.CompletableFuture;
@Component
@RequiredArgsConstructor
@Slf4j
public class DispatcherKafkaProducer {
private final KafkaTemplate kafkaTemplate;
public void send(String topic, String message) {
CompletableFuture future = kafkaTemplate.send(topic, message);
future.whenComplete((result, ex) -> {
throw new IllegalArgumentException("Exception occurred during sending message to Kafka topic.", ex);
});
}
}
< /code>
Сценарий проблемы < /h1>
Когда я запускаю приложение, я вижу, как запускается по продюсеру Kafka, но на моем кафке-кончителе не возвращает одно сообщение.2025-08-21T16:58:00.035+02:00 INFO 10328 --- [Scheduler] [ scheduling-1] o.a.k.clients.producer.ProducerConfig : ProducerConfig values:
acks = -1
auto.include.jmx.reporter = true
batch.size = 16384
bootstrap.servers = [localhost:9092]
buffer.memory = 33554432
client.dns.lookup = use_all_dns_ips
client.id = Scheduler-producer-1
compression.gzip.level = -1
compression.lz4.level = 9
compression.type = none
compression.zstd.level = 3
connections.max.idle.ms = 540000
delivery.timeout.ms = 120000
enable.idempotence = true
enable.metrics.push = false
interceptor.classes = []
key.serializer = class org.apache.kafka.common.serialization.StringSerializer
linger.ms = 0
max.block.ms = 60000
max.in.flight.requests.per.connection = 5
max.request.size = 1048576
metadata.max.age.ms = 300000
metadata.max.idle.ms = 300000
metadata.recovery.strategy = none
metric.reporters = []
metrics.num.samples = 2
metrics.recording.level = INFO
metrics.sample.window.ms = 30000
partitioner.adaptive.partitioning.enable = true
partitioner.availability.timeout.ms = 0
partitioner.class = null
partitioner.ignore.keys = false
receive.buffer.bytes = 32768
reconnect.backoff.max.ms = 1000
reconnect.backoff.ms = 50
request.timeout.ms = 30000
retries = 2147483647
retry.backoff.max.ms = 1000
retry.backoff.ms = 100
sasl.client.callback.handler.class = null
sasl.jaas.config = null
sasl.kerberos.kinit.cmd = /usr/bin/kinit
sasl.kerberos.min.time.before.relogin = 60000
sasl.kerberos.service.name = null
sasl.kerberos.ticket.renew.jitter = 0.05
sasl.kerberos.ticket.renew.window.factor = 0.8
sasl.login.callback.handler.class = null
sasl.login.class = null
sasl.login.connect.timeout.ms = null
sasl.login.read.timeout.ms = null
sasl.login.refresh.buffer.seconds = 300
sasl.login.refresh.min.period.seconds = 60
sasl.login.refresh.window.factor = 0.8
sasl.login.refresh.window.jitter = 0.05
sasl.login.retry.backoff.max.ms = 10000
sasl.login.retry.backoff.ms = 100
sasl.mechanism = GSSAPI
sasl.oauthbearer.clock.skew.seconds = 30
sasl.oauthbearer.expected.audience = null
sasl.oauthbearer.expected.issuer = null
sasl.oauthbearer.header.urlencode = false
sasl.oauthbearer.jwks.endpoint.refresh.ms = 3600000
sasl.oauthbearer.jwks.endpoint.retry.backoff.max.ms = 10000
sasl.oauthbearer.jwks.endpoint.retry.backoff.ms = 100
sasl.oauthbearer.jwks.endpoint.url = null
sasl.oauthbearer.scope.claim.name = scope
sasl.oauthbearer.sub.claim.name = sub
sasl.oauthbearer.token.endpoint.url = null
security.protocol = PLAINTEXT
security.providers = null
send.buffer.bytes = 131072
socket.connection.setup.timeout.max.ms = 30000
socket.connection.setup.timeout.ms = 10000
ssl.cipher.suites = null
ssl.enabled.protocols = [TLSv1.2, TLSv1.3]
ssl.endpoint.identification.algorithm = https
ssl.engine.factory.class = null
ssl.key.password = null
ssl.keymanager.algorithm = SunX509
ssl.keystore.certificate.chain = null
ssl.keystore.key = null
ssl.keystore.location = null
ssl.keystore.password = null
ssl.keystore.type = JKS
ssl.protocol = TLSv1.3
ssl.provider = null
ssl.secure.random.implementation = null
ssl.trustmanager.algorithm = PKIX
ssl.truststore.certificates = null
ssl.truststore.location = null
ssl.truststore.password = null
ssl.truststore.type = JKS
transaction.timeout.ms = 60000
transactional.id = null
value.serializer = class org.apache.kafka.common.serialization.StringSerializer
< /code>
И это журнал, когда мое приложение пытается отправить сообщение. < /p>
2025-08-21T16:59:00.041+02:00 DEBUG 10328 --- [toekeloo] [uler-producer-1] o.a.k.c.p.internals.RecordAccumulator : [Producer clientId=toekeloo-producer-1] Assigned producerId 2014 and producerEpoch 0 to batch with base sequence 1 being sent to partition SCH-dispatcher-0
2025-08-21T16:59:00.041+02:00 DEBUG 10328 --- [toekeloo] [uler-producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=toekeloo-producer-1] Sending PRODUCE request with header RequestHeader(apiKey=PRODUCE, apiVersion=11, clientId=toekeloo-producer-1, correlationId=5, headerVersion=2) and timeout 30000 to node 1: {acks=-1,timeout=30000,partitionSizes=[SCH-dispatcher-0=227]}
2025-08-21T16:59:00.042+02:00 DEBUG 10328 --- [toekeloo] [uler-producer-1] org.apache.kafka.clients.NetworkClient : [Producer clientId=toekeloo-producer-1] Received PRODUCE response from node 1 for request with header RequestHeader(apiKey=PRODUCE, apiVersion=11, clientId=toekeloo-producer-1, correlationId=5, headerVersion=2): ProduceResponseData(responses=[TopicProduceResponse(name='SCH-dispatcher', partitionResponses=[PartitionProduceResponse(index=0, errorCode=0, baseOffset=3292, logAppendTimeMs=-1, logStartOffset=2618, recordErrors=[], errorMessage=null, currentLeader=LeaderIdAndEpoch(leaderId=-1, leaderEpoch=-1))])], throttleTimeMs=0, nodeEndpoints=[])
< /code>
Я слишком незнаком, чтобы самостоятельно обнаружить основную причину проблемы. Итак, я ищу совета, почему мои сообщения, созданные кафкой, не получены брокером или почему брокер не потребляет сообщения.
Подробнее здесь: https://stackoverflow.com/questions/797 ... y-messages