# 운영 프로파일에 KafkaProfileValidator 가 거는 두 요구
KafkaProfileValidator.java:51 if (profile.production() && !profile.tlsEnabled()) {
KafkaProfileValidator.java:52 throw new IllegalArgumentException(
KafkaProfileValidator.java:53 "a production Kafka connection requires TLS: " + profile.broker());
KafkaProfileValidator.java:54 }
KafkaProfileValidator.java:55 if (profile.production() && !profile.authenticationEnabled()) {
KafkaProfileValidator.java:56 throw new IllegalArgumentException(
KafkaProfileValidator.java:57 "a production Kafka connection requires broker authentication: " + profile.broker());
KafkaProfileValidator.java:58 }
# 그 검증기가 기동에서 도는 경로와, 그 거부를 지키는 테스트
main · KafkaMessagingAutoConfiguration.java:55 kafkaProfileStartupValidation(
main · StartupProfileValidation.java:19 *
Validation runs as {@code afterPropertiesSet} rather than on an application event, so the
main · StartupProfileValidation.java:31 final class StartupProfileValidation implements InitializingBean {
main · StartupProfileValidation.java:43 public void afterPropertiesSet() {
test · MessagingConfigurationBindingTest.java:235 void aProductionKafkaBrokerWithoutTransportSecurityFailsStartup() {
# 두 조립부를 켜는 프로퍼티
main · MessagingBridgeRootAutoConfiguration.java:21 @ConditionalOnProperty(prefix = "app.messaging", name = "enabled", havingValue = "true")
main · MessagingAdminAutoConfiguration.java:27 @ConditionalOnProperty(prefix = "app.messaging.admin", name = "enabled", havingValue = "true")
main · MessagingPlatformRootAutoConfiguration.java:28 @ConditionalOnProperty(prefix = MessagingSettings.PREFIX, name = "enabled", havingValue = "true")
main · MessagingProviderSelection.java:48 "kafka", KafkaMessagingAutoConfiguration.class,
main · KafkaSenderConfig.java:38 @ConditionalOnProperty(name = "app.messaging.broker", havingValue = "kafka")
출하 기본값 : application.yml:949: enabled: ${APP_MESSAGING_ENABLED:false}
# 프로덕션에서 KafkaProducer 를 만드는 두 자리 — 설정 맵 원문 그대로
KafkaMessagingAutoConfiguration.java:128 @Bean(destroyMethod = "close")
KafkaMessagingAutoConfiguration.java:129 @ConditionalOnMissingBean
KafkaMessagingAutoConfiguration.java:139 java.util.Map config = new java.util.HashMap<>();
KafkaMessagingAutoConfiguration.java:140 config.put(
KafkaMessagingAutoConfiguration.java:141 org.apache.kafka.clients.producer.ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
KafkaMessagingAutoConfiguration.java:142 bootstrapServers);
KafkaMessagingAutoConfiguration.java:143 config.put(
KafkaMessagingAutoConfiguration.java:144 org.apache.kafka.clients.producer.ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
KafkaMessagingAutoConfiguration.java:145 org.apache.kafka.common.serialization.ByteArraySerializer.class);
KafkaMessagingAutoConfiguration.java:146 config.put(
KafkaMessagingAutoConfiguration.java:147 org.apache.kafka.clients.producer.ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
KafkaMessagingAutoConfiguration.java:148 org.apache.kafka.common.serialization.ByteArraySerializer.class);
KafkaMessagingAutoConfiguration.java:149 // Durability is not a default worth inheriting: acks=1 loses an accepted publish to a leader
KafkaMessagingAutoConfiguration.java:150 // failover, which is exactly the outcome an outbox exists to prevent.
KafkaMessagingAutoConfiguration.java:151 config.put(org.apache.kafka.clients.producer.ProducerConfig.ACKS_CONFIG, "all");
KafkaMessagingAutoConfiguration.java:152 config.put(org.apache.kafka.clients.producer.ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
KafkaMessagingAutoConfiguration.java:153 return new org.apache.kafka.clients.producer.KafkaProducer<>(config);
KafkaSenderConfig.java:61 @Bean(name = "kafkaSeamProducer", destroyMethod = "close")
KafkaSenderConfig.java:62 @ConditionalOnMissingBean(name = "kafkaSeamProducer")
KafkaSenderConfig.java:64 Map config = new HashMap<>();
KafkaSenderConfig.java:65 config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, String.join(",", settings.brokers()));
KafkaSenderConfig.java:66 config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
KafkaSenderConfig.java:67 config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
KafkaSenderConfig.java:68 config.put(ProducerConfig.ACKS_CONFIG, "all");
KafkaSenderConfig.java:69 config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
KafkaSenderConfig.java:77 config.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 10_000);
KafkaSenderConfig.java:78 config.put(ProducerConfig.LINGER_MS_CONFIG, 0);
KafkaSenderConfig.java:79 config.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 15_000);
KafkaSenderConfig.java:80 config.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 30_000);
KafkaSenderConfig.java:81 return new KafkaProducer<>(config);
# security.protocol 을 넣는 프로덕션 코드 — 상수명과 리터럴 양쪽으로
main · KafkaSecurityConfigurer.java:32 public static final String SECURITY_PROTOCOL = "security.protocol";
main · KafkaSecurityConfigurer.java:80 properties.put(SECURITY_PROTOCOL, securityProtocol(profile, credential));
# KafkaSecurityConfigurer.configure 가 자격 종류마다 하는 일
KafkaSecurityConfigurer.java:90 case BrokerCredentialProfile.SaslScram scram -> {
KafkaSecurityConfigurer.java:92 properties.put(SASL_MECHANISM, "SCRAM-SHA-512");
KafkaSecurityConfigurer.java:93 properties.put(SASL_JAAS_CONFIG, scramJaas(scram.credentialId(), resolved));
KafkaSecurityConfigurer.java:95 case BrokerCredentialProfile.OAuth2 oauth -> {
KafkaSecurityConfigurer.java:99 throw new dev.caskeleton.messaging.api.error.MessagingConfigurationException(
KafkaSecurityConfigurer.java:106 case BrokerCredentialProfile.MutualTls ignored -> properties.put(SASL_MECHANISM, "NONE");
KafkaSecurityConfigurer.java:107 case BrokerCredentialProfile.UsernamePassword plain -> {
KafkaSecurityConfigurer.java:109 properties.put(SASL_MECHANISM, "PLAIN");
KafkaSecurityConfigurer.java:110 properties.put(SASL_JAAS_CONFIG, plainJaas(plain.credentialId(), resolved));
KafkaSecurityConfigurer.java:112 case BrokerCredentialProfile.Nkey ignored ->
KafkaSecurityConfigurer.java:113 throw new IllegalArgumentException(
# 그 클래스가 빈이 되기까지의 조건 사슬
KafkaMessagingAutoConfiguration.java:105 @Bean
KafkaMessagingAutoConfiguration.java:106 @ConditionalOnMissingBean
KafkaMessagingAutoConfiguration.java:107 @org.springframework.boot.autoconfigure.condition.ConditionalOnBean(
KafkaMessagingAutoConfiguration.java:108 CredentialRuntimeRegistry.class)
KafkaMessagingAutoConfiguration.java:109 public KafkaSecurityConfigurer kafkaSecurityConfigurer(
KafkaMessagingAutoConfiguration.java:110 CredentialRuntimeRegistry credentials, BrokerTlsPolicy tlsPolicy) {
KafkaMessagingAutoConfiguration.java:111 return new KafkaSecurityConfigurer(credentials, tlsPolicy);
KafkaMessagingAutoConfiguration.java:112 }
main · MessagingCoreAutoConfiguration.java:315 @ConditionalOnBean(CredentialProvider.class)
main · MessagingCoreAutoConfiguration.java:317 public CredentialRuntimeRegistry credentialRuntimeRegistry(CredentialProvider provider) {
CredentialProvider 를 구현하는 main 클래스 : 0 건
그 인터페이스를 구현하는 자리 전부 :
test · CredentialRuntimeRegistryTest.java:21 private static final class CountingProvider implements CredentialProvider {
# 그 빈을 파라미터나 필드로 받는 코드 (빈을 만드는 팩토리 자신은 뺀다)
main 에서 0 건
getBean · ObjectProvider · 빈 이름 문자열로 가져가는 자리 :
0 건
# 조립을 확인하는 테스트가 단언하는 것과, 실 브로커 시험이 붙는 컨테이너
test · MessagingStarterOffContractTest.java:118 void selectingKafkaAssemblesOnlyKafka() {
test · MessagingStarterOffContractTest.java:172 void aSelectedTransportAssemblesAPublisher() {
test · MessagingStarterOffContractTest.java:127 assertThat(context).hasSingleBean(KafkaMessagingAutoConfiguration.class);
test · MessagingStarterOffContractTest.java:130 .doesNotHaveBean(RabbitMessagingAutoConfiguration.class);
test · MessagingStarterOffContractTest.java:189 assertThat(context).hasSingleBean(MessagePublisher.class);
test · MessagingLiveRoundTripQualificationTest.java:66 @Container static final KafkaContainer KAFKA = new KafkaContainer("apache/kafka:4.1.0");