Spring발행일 2025. 7. 1.원본 https://blog.naver.com/jword_/223918269974 ↗

Kafka 사용하기 (Producer, Consumer, DLQ)

Kafka 사용하기 (Producer, Consumer, DLQ) — #개발자의도구들 #AI스쿨 #Kafka #KafakProducer #카프카 AI스쿨 msa기반 java 백엔드 코스 중에 공...

#Spring#Naver Blog

#개발자의도구들 #AI스쿨 #Kafka #KafakProducer #카프카

​

AI스쿨 msa기반 java 백엔드 코스 중에 공부한 내용을 작성하였습니다

\* 현재 Spring x 프로젝트를 진행하고 있습니다. 전체 목차를 보시려면 여기를 눌러주세요.

\* Devops stack: 여기

\* 리눅스 기본기: 여기

Spring Boot에서 Kafka를 사용하기

Noti 기능에 주로 Kafka서비스를 사용할 수 있습니다. 이전글을 통해 Noti 설계 방법을 배워봤으니 이제 안정성을 더하는 Kafka를 응용해봅시다.

왜 Kafka를 써야하는가?

현재 X 플랫폼이 MSA 구조로 나눠져 있습니다.

이미지

Noti의 경우 NotiServer 내에서 데이터를 가지고 있는 구조입니다. Front Server에서 특정 Event가 발생하면 Front Server는 Noti Server에 "이런 이런 메세지를 저장해줘" 라고 데이터를 보내게 될 것 입니다.

​

이미지

이 구조에서 FrontServer와 NotiServer의 연결이 끊어진다면 어떻게 될까요? 프론트에서 일어난 특정 이벤트에 대한 Noti를 NotiServer가 전송받지 못해 결국 사용자에게 Noti를 전송하지 못합니다.

근본적으로 이런 메세지 전송 오류 문제를 해결하기 위해 다양한 솔루션이 개발되었으며, Kafka는 대표적인 대용량 메세지 처리 서비스가 되겠습니다.

​

이미지

그림으로 나타내면 아런 구조가 되겠네요. 여기서 Front Server는 Producer의 역할을 Noti Server는 Consumer 역할을 수행하며, Kafka 서비스는 이를 관리하는 Controller(Broker) 서비스를 제공합니다.

​

우리가 해야할 일은 아래와 같습니다.

text 코드 예제
                                    0. kafka를 실행하기

1. Spring Boot Project와 Kafka를 연동하기

2. Front Server를 Producer로 등록하기

3. Noti Server를 Consumer로 등록하기

​

구현 코드

📌 0. Kafka 를 실행하기

​

저는 Linux에 설치를 진행하였는데, 이 부분은 분량이 조금 길어질 것 같아 따로 정리하도록 하겠습니다.

​

​

📌 1. Spring boot 와 Kafka를 연동하자.

kotlin 코드 예제
                                    implementation("org.apache.kafka:kafka-clients:3.9.1")

build.gradle에 kafka-client 모듈을 추가합니다. 이는 Spring Boot 프로젝트가 Kafka와 직접 연결가능한 Client 모듈을 제공합니다.

​

kotlin 코드 예제
                                    @Configuration
class KafkaProducerConfig {

    @Bean
    fun kafkaProducer(
        @Value("\${spring.kafka.bootstrap-servers}") kafkaBootstrapServers: String,
        @Value("\${spring.kafka.producer.key-serializer}") keySerializer: String,
        @Value("\${spring.kafka.producer.value-serializer}") valueSerializer: String
    ): KafkaProducer<String, String> {
        val props = Properties().apply {
            put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers)
            put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, keySerializer)
            put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, valueSerializer)
        }
        return KafkaProducer<String, String>(props)
    }
}

Configuration을 등록합니다. 저의 경우 kafka 전용 yaml을 생성하고 뷸러와서 사용하였습니다. 핵심 적인 부분은 다음과 같습니다.

​

  1. Kafka가 구동중인 서버 주소
  1. Kafka와 소통하기 위한 key-value Serializer

​

yaml 코드 예제
                                    spring:
  kafka:
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    bootstrap-servers: << Kafka 서버 주소 >>

key-value Serializer는 Front Server에서 생성된 객체를 Kafka 브로커로 전송하기 위해 Byte\[\]로 변환하는 역할을 수행합니다. Kafka는 이렇게 변환된 값을 이해하여 topic과 value를 저장할 수 있습니다.

​


​

📌 2. Spring boot Producer 제작

kotlin 코드 예제
                                    @Service
class KafkaProducerService(
    private val kafkaProducer: KafkaProducer<String, String>,
    private val objectMapper: ObjectMapper
) {
    private val logger = KotlinLogging.logger {}

    fun sendNoti(request: NotificationSaveRequest) {
        val topicName = "noti" // 예: Consumer가 구독 중인 동일 토픽
        try {
            // DTO -> JSON
            val jsonValue = objectMapper.writeValueAsString(request)

            // ProducerRecord 생성 후 send
            val record = ProducerRecord<String, String>(topicName, jsonValue)
            kafkaProducer.send(record) { metadata, exception ->
                if (exception == null) {
                    logger.info {
                        "Sent message to $topicName [offset=${metadata.offset()}, " +
                                "partition=${metadata.partition()}]" }
                } else {
                    logger.error { "Failed to send message: ${exception.message}" }
                }
            }
        } catch (e: Exception) {
            logger.error { "sendMessage error: ${e.message}" }
        }
    }
}

Front Server 객체를 JSON으로 변환하기 위해 Object Mapper를 사용하였습니다.

text 코드 예제
                                    ✍️  Spring 전용 Kafka module을 사용하시면
Object Mapper 없이 JsonSerializer를 직접 사용할 수 있습니다.

중요 핵심사항은 Kafka에 어떤 데이터를 어떻게 줄 것인지 잘 정의하는 것입니다. 저는 topic을 Noti로 등록하고 NotificadtionSaveRequst 객체를 value로 전달하고 있습니다.

​


​

📌 3. Spring boot Consumer 고려사항

\* 구현 코드는 생략하였습니다. Error 사항 정리 후 다시 돌아와서 작성하겟습니다.

​

Consumer 부분에서 꽤나 애를 먹었었습니다. 이유는 Consumer의 경우 고려해야하는 경우의수가 발생하기 때문입니다.

​

자, 다시한번 그림을 봅시다.

이미지

우선 Front에서 Kafka로 올바르게 Noti가 전송되었다고 해보겠습니다. 이후 Kafka는 해당 메세지를 가지고 있는 상태입니다.

​

이때 Kafka의 Consumer인 Noti Server는 특정 Topiv에 대한 Noti를 불러오고 싶어 합니다. 하지만 실패했네요! 그럼 이제 어떻게 해야할까요?

​

어떻게 해야하는지 알아보기 이전에 한발짝 물러서서 🤔 원인을 먼저 분석해봅시다. 원인을 알면 어떻게 해야할지 길이 보이기 때문입니다.

​

text 코드 예제
                                    🤔 왜 실패하는가?

* Kafka는 정상이라고 가정합니다. 우리가 대응 해야 하는 것은 Consumer 코드입니다.

1. Noti Server와 Kafka의 네트워크 오류
- 거의 발생하지 않지만 일시적으로 네트워크 오류로 인해 패킷이 손실 될 수 있습니다.

2. Noti Server의 DB오류
- Noti Server의 DB를 일시적으로 사용할 수 없다면, 해당 Noti는 저장되지 않습니다.
 저장되지 않은 Noti는 Client에게 제공될 수 없습니다.

3. Packet Parsing 오류
- Kafka Value로 저장된 데이터는 Byte[] 입니다. 이를 JSON으로 다시 Parsing 하기 위해
Noti Server에서 Object Mapper사용하여 변환을 시도할 것 입니다.

- 하지만 이때 잘못된 필드명이나 데이터 형식이 들어있다면, 오류가 발생할 수 있습니다.

Producer와 달리 Consumer는 더 복잡한 비즈니스 로직이 필요합니다. 이는 회사 정책에 따라 어떻게 오류를 관리해야할지 정해지지만, 일반적으로 대응 가능한 방법들을 공유하고자 합니다.

​

​

⚠️ 에러 대응 1. 일시적 네트워크 오류

​

일시적으로 네트워크가 오류나여 제대로 데이터를 받아보지 못하였다면, 다시 요청하여 받으면 됩니다. Kafka는 기본적으로 commit되지 않으면 consume되지 않았다고 보며, commit 되더라도 일정 기간동안 해당 데이터를 삭제하지 않습니다.

text 코드 예제
                                    ✍️ solution
몇 번 retry 해보고 그래도 계속 안되면 더 이상 consume 진행에 의미가 없기에 나중에 다시 시도

⚠️ 에러 대응 2. DB 오류

DB에 문제가 생겨 Noti를 올바르게 저장할 수 없다면 Noti Service 자체 의미가 없습니다. 이 경우 DB가 복구되기까지 기다렸다가 다시 시도하는 것이 좋습니다.

text 코드 예제
                                    ✍️ solution
여기도 몇 번 retry 해보고 계속 안되면 더 이상 consume의 의미가 없으므로 consume 종료

⚠️ 에러 대응 3. Parsing 오류

사실 99.9%는 해당 Parsing오류라고 볼 수 있습니다. Parsing 오류는 개발자들 사이에서의 실수에서 비롯되기 때문입니다. 특히나 JSON 데이터를 관리할 때 키 값이나, 데이터 형식을 제대로 안적는다거나, 혹은 특정 필드를 빼놓고 적는다거나 등 여러가지 문제가 발생할 수 있습니다.

​

이 경우에는 여러가지 솔루션들이 있습니다.

text 코드 예제
                                    ✍️ solutions

Json Parsing Error가 발생하면 retry 없이 아래 3개의 방법 중 하나 이상을 선택한다.

1. log에 남긴다

2. 다른 queue에 남긴다

3. 다른 table에 남긴다

 데이터를 남긴 후 해당 메세지에 대한 commit을 진행한다. 치명적인 오류가 아니기 때문에 consume을
종료할 이유는 없다.

솔루션의 공통점은 흔적을 남기는 것 입니다. 어떤 부분이 오류인지 정확히 파악해야하고, 또한 소비된 메세지를 그냥 버리는 것이 아니라, 차후 수정을 통해 DB에 제대로 저장 시킬 수 있기 때문입니다.

​

🚀 DLQ

  • *Dead Letter Queue의 약자로 오류난 메세지를 따로 관리하기 위해 Queue에 저장해 둡니다. Kafka에서는 해당 Queue를 새로운 Topic으로 발행할 수 있습니다.

​

구현코드

kotlin 코드 예제
                                    @Service
class KafkaConsumerRunner(
    private val kafkaConsumer: KafkaConsumer<String, String>,
    private val objectMapper: ObjectMapper,
    private val notificationRepository: NotificationRepository
) {
    private val logger = KotlinLogging.logger {}
    private val kafkaDlqLogger = LoggerFactory.getLogger("com.example.notiApiServer.DLQ")

    @Volatile
    private var consumerRunning = true

    fun poll(vararg args: String?) {
        logger.info {"start consumer runner poll ()"}
        val topicName = "noti"
        kafkaConsumer.subscribe(listOf(topicName))

        logger.info {"consume starts"}

        // Graceful shutdown을 위한 shutdown hook 등록
        Runtime.getRuntime().addShutdownHook(Thread {
            println("Shutdown initiated.")
            turnOff()
        })

        while (consumerRunning) {
            val records = kafkaConsumer.poll(Duration.ofSeconds(3000))

            for (record in records) {
                logger.info {"record : ${record.topic()}, ${record.offset()}, ${record.value()}"}
                try {
                    val notiSaveRequest = recordToNotiSaveRequest(record)
                    // save in DB
                    val notification = Notification(
                        id = null,
                        publisherId = notiSaveRequest.publisherId,
                        receiverId = notiSaveRequest.receiverId,
                        notificationType = notiSaveRequest.notificationType,
                        targetBoardId = notiSaveRequest.targetBoardId,
                        createdAt = null
                    )
                    // save db
                    saveNoti(notification)

                    val topicPartition = TopicPartition(record.topic(), record.partition())
                    val offsetAndMetadata = OffsetAndMetadata(record.offset() + 1)
                    kafkaConsumer.commitSync(mapOf(topicPartition to offsetAndMetadata))

                    // commit offset !
                    logger.info { "Record with offset ${record.offset()} processed and committed." }
                } catch (e: JsonProcessingException) {
                    // skip this record
                    continue
                } catch (e: DataAccessException) {
                    logger.error {"DB error occurred!! at offset: ${record.offset()}"}
                }
            }
        }
    }

    private fun recordToNotiSaveRequest(record: ConsumerRecord<String, String>): NotificationSaveRequest {
        val topic = record.topic()

        val partition = record.partition()
        val offset = record.offset()
        val json = record.value()
        val key = Triple(topic, partition, offset)
        try {
            return objectMapper.readValue<NotificationSaveRequest>(json)
        } catch (e: JsonProcessingException) {
            logger.error {"toNotiSaveRequest parsing error"}
            // save log
            val errorDto = KafkaDlqRecord(
                topic = topic,
                value = json,
                offset = offset,
                revive = false,
                reviveDate = null
            )
            logKafkaDlqError(errorDto)
            // offset -> +1 enforce
            val topicPartition = TopicPartition(topic, partition)
            val offsetData = OffsetAndMetadata(offset + 1) // offset commit to next
            kafkaConsumer.commitSync(mapOf(topicPartition to offsetData))
            throw e
        }
    }

    private fun logKafkaDlqError(dlqRecord: KafkaDlqRecord) {
        try {
            MDC.put("topic", dlqRecord.topic)
            MDC.put("offset", dlqRecord.offset.toString())
            MDC.put("value", dlqRecord.value)
            MDC.put("revive", dlqRecord.revive.toString())
            MDC.put("reviveDate", dlqRecord.reviveDate.toString())
            kafkaDlqLogger.error("Kafka DLQ error occurred: {}", dlqRecord)
        } finally {
            MDC.clear()
        }
    }

    private fun saveNoti(noti: Notification) {
        try {
            notificationRepository.save(noti)
        } catch (e: DataAccessException) {
            throw e
        }
    }

    private fun turnOff() {
        consumerRunning = false
        kafkaConsumer.close()
    }

    private fun turnOn() {
        // scheduling?
    }
}

풀리지 않는 의문들

🤔 Kafka 자체가 Down 되면 각 서버는 어떻게 동작하는가?

​

기본적으로 Kafka Client의 Producer, Consumer는 서버가 살아있는 동안 끊임없이 연결을 시도합니다. Kafka가 Down되면 새로운 Noti 메세지를 전달할 수 없지만, 기존의 Noti는 받아볼 수 있습니다. 이는 Front와 Noti Server가 Kafka이외에 따로 연결을 시도하기 때문입니다.

​

새로 발행된 Noti를 관리하는 전략은 다양하게 있습니다. 하지만 근본적인 전략은 아래와 같습니다.

​

text 코드 예제
                                    Producer에서 Noti를 가지고 있다가 Kafka가 복구되면 다시 발행을 시도한다.

위에서 언급한 것 처럼 Noti를 저장하는 방법은 무궁무진하기 때문에 회사 정책에 맞게 잘 관리하시면 되겠습니다!