https://docs.spring.io/spring-cloud-stream/docs/current/reference/html/spring-cloud-stream.html#spring-cloud-stream-overview-introducing
메세징 서비스와 분산캐싱
StreamBridge 사용하기
스프링 클라우드 스트림은 RabbitMQ나 아파치 카프카 등 여러 메세지 플랫폼과의 바인딩 부분을 추상화하여 제공하는데 그로 인해 개발자는 메세지 플랫폼에 구애받지 않고 개발을 할 수 있다.
또한 간단한 어노테이션을 활용한 개발을 통해 publisher와 consumer 구현을 간편하게 해준다는 장점이 있다.
implementation(Spring.cloud.stream.binderKafka)
implementation("org.springframework.cloud:spring-cloud-starter-stream-kafka")
implementation(Spring.cloud.stream.stream)
cloud:
function:
definition: kafkaCloudSender #cloud stream을 사용할 클래스
stream:
binders: #바인더 생성
kafka-binder: #이름 설정
type: kafka
environment:
spring.cloud.stream.kafka.binder.brokers: ${KAFKA_BROKERS} #카프카 브로커설정
bindings:
kafkaCloudSender-out-0: #바인딩 이름
binder: kafka-binder #위에서 정의한 바인더 이름
destination: event.search # topic
contentType: application/json
message 파라미터의 Event 타입은 사내에서 정의된 타입이라 이부분은 각자 구현하면 됩니당~~
class KafkaCloudSender(
val streamBridge StreamBridge // StreamBridge 주입
){
//binding은 yml의 bindings에서 정의한 바인딩 이름
@Async
fun send(binding: String, message: Event){
streamBridge.send(binding, MessageBuilder
.withPayload(message)
.setHeader(KafkaHeaders.MESSAGE_KEY, UUID.randomUUID().toString())
.build())
}
}
override suspend fun updateProductDoc(productId: String, version: String) {
val binding = "KafkaCloudSender-out-0"
val event = generateSchema(productId = productId).toEvent(version = version)
kafkaCloudSender.send(binding = binding, message = event)
}
필요한 로직에서 위에서 만든 KafkaCloudSender의 send를 사용해서 publish 해주면 끝!