grep

Backend

Kafka Message Conversion (Serial/Deserializing)

Ashe여기어때

2025년 2월 6일

원문에서 보기 ↗

안녕하세요. 여기어때컴퍼니 주문결제개발팀 애쉬입니다.

프로젝트 진행 도중 ‘Kafka Message를 Consume할 때 Json Value를 Parsing하는 과정이 좀 이상한데?’에서 시작한 짧은 이야기를 들려 드릴까 합니다. 한 편의 수필이라 생각하고 봐주시면 감사하겠습니다.

Project versions

java openjdk 17
spring boot 2.7.18
apache kafka 3.1.2
spring kafka 2.8.11

Json과 함께하는 Kafka 여행

Kafka 메세지 변환 관련 설정을 뜯어보게 된 이유?

기존에 Kafka Message를 Consume하는 우리 서비스들은 대부분 메세지를 Key 없이 Value를 String으로 받아서 ObjectMapper를 활용해서 변환을 하고 있었어.

private ConsumerFactory<String, String> setConsumerFactory(final KafkaProperties.ConsumerConfig consumerConfig) {
    Map<String, Object> configProps = new HashMap<>();
    ... 중략
    configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    ... 중략
    return new DefaultKafkaConsumerFactory<>(configProps);
}
@KafkaListener(containerFactory = "testKafkaListenerContainerFactory", topics = "${spring.kafka.consumer.test.topic}")
public void consumeTest(ConsumerRecord<String, String> consumerRecord) {
    try {
        TestDTO testDTO = objectMapper.readValue(consumerRecord.value(), TestDTO.class);
        // do something with DTO
    } catch (Exception e) {
        // do something with exception
    }
}

그래서 아 Producer 쪽에서 Json String으로 모델을 Converting해서 보내고 있나 보다! 하고 Producer 서비스를 열어보았는데 웬걸 Message Value를 JsonSerializer를 통해서 (objectMapper를 통한 jsonString 생성 후 전달이 아닌) 직변환해서 사용 중이었더라고?

@Bean
public ProducerFactory<String, Object> producerFactory() {
    KafkaProperties.ProducerCommonConfig config = kafkaProperties.getProducer().getCommonConfig();
    Map<String, Object> configProps = new HashMap<>();
    ... 중략
    configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); // ???
    ... 중략
    return new DefaultKafkaProducerFactory<>(configProps);
}

그럼 Consumer 설정에서도 JsonDeserializer로 바로 변환해서 받을 수 있어야 하는거 아닌가?

이 질문에 답을 찾기 위해 참 많이도 돌아다녔다 😇 그렇게 답을 찾기 위한 여정을 떠났지… (함께 여행한 CheckedOutMessage 야 고마워)

(ps. journey는 첫 개발을 시작할 때 팀장님의 닉네임이었어, 이제야 좀 이해가 되는 닉네임이네 ㅋㅋ)

Journey Season 1 (JsonDeserializer 설정)

JsonDeserializer 설정하면 끝 아님? < ㅋㅋ 그게 됐으면 이 글을 쓰고 있겠냐

@Bean
public ConcurrentKafkaListenerContainerFactory<String, CheckedOutMessage> checkedOutKafkaListenerContainerFactory() {
    // Consumer 설정
    Map<String, Object> consumerConfigProps = getCommonConsumerConfigProps(); // 기본 서버정보 등

    // Container Factory 설정
    ConcurrentKafkaListenerContainerFactory<String, CheckedOutMessage> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigProps, 
    new StringDeserializer(), 
    new JsonDeserializer<>(CheckedOutMessage.class))); // 일케 넣음 알아서 변환 해주겠지 라고 생각한 내가 개발자가 맞나 싶다
    return factory;
}

제일 처음 시도한 내용은 DefaultKafkaConsumerFactory 를 생성할 때 properties, keyDeserializer, valueDeserializer를 설정하는 데 여기를 StringDeserializer > JsonDeserializer로 변환하는 거였어.

Mapping Model을 전달하면 알아서 잘 변환 하겠지? 히히? 하고 신나게 테스트를 돌렸는데

?

1차 오류 (오류가 무한대로 난다. 이건 뒤에서 다시 볼게)

Consumer exception
java.lang.IllegalArgumentException: The class 'test.order.kafka.message.CheckedOutMessage' is not in the trusted packages: 
[java.util, java.lang, test.orderWorker.consumer.message, test.orderWorker.consumer.message.*]. 
If you believe this class is safe to deserialize, please provide its name. 
If the serialization is only done by a trusted source, you can also enable trust all (*).

아 니가 발행처의 메세지를 신뢰할 수 없다고? 그래 그러면 시키는 대로 전체 허용 (“*”) 해서 해 볼게?

JsonDeserializer<CheckedOutMessage> jsonDeserializer = new JsonDeserializer<>(CheckedOutMessage.class);
jsonDeserializer.addTrustedPackages("*");
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigProps, new StringDeserializer(), jsonDeserializer));

한다~?

? 아니 시키는 대로 했잖아

2차 오류

Consumer exception
java.lang.ClassNotFoundException: test.order.kafka.message.CheckedOutMessage
    ... 중략
 at org.springframework.util.ClassUtils.forName(ClassUtils.java:291)
 at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.getClassIdType(DefaultJackson2JavaTypeMapper.java:138)
 ... 중략
Wrapped by: org.springframework.messaging.converter.MessageConversionException: failed to resolve class name. 
Class not found [test.order.kafka.message.CheckedOutMessage]; nested exception is 
java.lang.ClassNotFoundException: test.order.kafka.message.CheckedOutMessage

아 그렇지 orderWorker 프로젝트에서 order에서 발행한 메세지 DTO 경로에 있는 파일이 있을 리가 없잖아.

아니 그러면 orderWorker의 test.orderWorker.kafka.message.CheckedOutMessage 와 발행처 order의 test.order.kafka.message.CheckedOutMessage 메세지가 같은 거라 mapping 해서 변환해 달라고 어떻게 deserializer한테 전달할 수 있지?

Journey Season 2 (Json parsing with mapping types)

Serialization, Deserialization, and Message Conversion :: Spring Kafka

역시 공식 문서가 짱이야! (GPT 너 답 똑바로 못하더라? 쓰읍) 시키는 대로 하니까 정말 잘 동작 하잖아?

consumerConfigProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
consumerConfigProps.put(JsonDeserializer.TYPE_MAPPINGS, "checkedOut:test.orderWorker.kafka.message.CheckedOutMessage");
consumerConfigProps.put(JsonDeserializer.TRUSTED_PACKAGES, "test.*");

JsonDeserializer를 설정하고 JsonDeserializer.TYPE_MAPPINGS 를 설정하는데 token:className 형태로 지정해 주면 돼

token : 별칭(producer와 consumer간에 약속된 별칭으로 지정 필수)

className :

producer쪽은 발행하는 메세지의 fullPath(ex. test.order.kafka.message.CheckedOutMessage),

consumer쪽은 매핑하는 메세지의 fullPath(ex. test.orderWorker.kafka.message.CheckedOutMessage)

결국 이 방식은 producer 쪽과 지정된 topic에서 지정된 message를 같은 token(별칭)을 사용하여 이게 둘이 같은 Message 야!를 알려주고 변환하는 방식이었어. 이 방식의 장점은 신뢰할 수 있는 메세지만 수신하기 때문에 white list 관리뿐만 아니라 변경, 메세지 오염 등을 수신 시점에 변환할 때 발견할 수 있다는 점이야.

but, string으로 직접 className을 지정하는 만큼 human error가 많이 발생할 수 있는 방법이고, producer 쪽에도 TYPE_MAPPING 을 지정해야 하는 관리 포인트가 생겨나게 되더라고(이쪽으로 떠났던 여행은 너무 산으로 가는 내용이라 따로 적지 않았어.. 하하.. 😇)

그래서 이 방식으로 내부에 공유를 했는데 우리 팀의 에이스 @파오리Paori(박지형) / 주문결제개발팀에게 휴먼 에러와 관리 포인트가 늘어나는 것에 비해 얻는 이득이 큰 것 같지 않다는 피드백을 받고 다시 여행을 떠나게 되었어😋 (죽여줘)

Journey Season 3 (Ignore type information in headers with useHeadersIfPresent)

위에 있는 공식 문서에서 조금 아래로 내려보면 objectMapper를 재정의해서 사용할 필요가 있다는 예시를 지나면

Starting with version 2.2, you can explicitly configure the deserializer to use the supplied target type and ignore type information in headers by using one of the overloaded constructors that have a boolean useHeadersIfPresent argument (which is true by default). The following example shows how to do so:

DefaultKafkaConsumerFactory<Integer, Cat1> cf = new DefaultKafkaConsumerFactory<>(props,
        new IntegerDeserializer(), new JsonDeserializer<>(Cat1.class, false));

이런 설명을 볼 수 있는데 처음 보면 이해가 잘되지 않았어 하지만 일단 시키는 대로 해봐야겠지?

factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigProps, 
new StringDeserializer(), 
new JsonDeserializer<>(CheckedOutMessage.class, false)));

Eureka!!!!!

자 잘 동작하니까 이제 마무리하고 어서 프로젝트 개발을 해야하지 않을까? 너무 여행을 멀리 떠나버려서일까? 이해를 못 하고 넘어가는 게 싫어서 소스를 좀 파고들어갔어 😢 (왜 그랬어..)

Journey Season 4 (JsonDeserializer 타입 추론 방식)

DefaultKafkaConsumerFactory의 생성자는 이렇게 생겼어

public DefaultKafkaConsumerFactory(Map<String, Object> configs, @Nullable Deserializer<K> keyDeserializer, @Nullable Deserializer<V> valueDeserializer) {
    this(configs, () -> {
        return keyDeserializer;
    }, () -> {
        return valueDeserializer;
    });

config 정보들과 key, value deserializer들을 설정할 수 있지. 우리 팀은 key를 지정해서 특정 파티션이 소비하도록 설정하지 않고 랜덤 분배하도록 되어 있기 때문에 value jsonDeserializer 에만 집중해 봤어.

public JsonDeserializer(@Nullable Class<? super T> targetType, boolean useHeadersIfPresent) {
    this(targetType, JacksonUtils.enhancedObjectMapper(), useHeadersIfPresent);
}

또 소스를 타고 가다 보면

private void initialize(@Nullable JavaType type, boolean useHeadersIfPresent) {
    this.targetType = type;
    this.useTypeHeaders = useHeadersIfPresent;
    Assert.isTrue(this.targetType != null || useHeadersIfPresent, "'targetType' cannot be null if 'useHeadersIfPresent' is false");
    if (this.targetType != null) {
        this.reader = this.objectMapper.readerFor(this.targetType);
    }
    this.addTargetPackageToTrusted();
    this.typeMapper.setTypePrecedence(useHeadersIfPresent ? TypePrecedence.TYPE_ID : TypePrecedence.INFERRED);
}

useHeadersIfPresent 여부에 따라 typePrecedence를 TYPE_ID or INFERRED로 설정하게 되어 있어. 이름에서도 알 수 있듯이 typeMapper 가 어떤 설정을 우선시할 것인지를 세팅하는데 기본은 (TYPE_ID) 그리고 useHeadersIfPresent가 false로 오게 되면 INFERRED(추론)으로 변경돼. 그럼 이 typeMapper는 무엇을 하는 녀석일까? 아래는 실제 메세지를 deserialize 할 때 호출되는 메소드야

public T deserialize(String topic, Headers headers, byte[] data) {
    if (data == null) {
        return null;
    } else {
        ObjectReader deserReader = null;
        JavaType javaType = null;
        if (this.typeResolver != null) {
            javaType = this.typeResolver.resolveType(topic, data, headers);
        }
        if (javaType == null && this.typeMapper.getTypePrecedence().equals(TypePrecedence.TYPE_ID)) {
            javaType = this.typeMapper.toJavaType(headers);
        }
        if (javaType != null) {
            deserReader = this.objectMapper.readerFor(javaType);
        }
        if (this.removeTypeHeaders) {
            this.typeMapper.removeHeaders(headers);
        }
        if (deserReader == null) {
            deserReader = this.reader;
        }
        Assert.state(deserReader != null, "No type information in headers and no default type provided");
        try {
            return deserReader.readValue(data);
        } catch (IOException var7) {
            IOException e = var7;
            throw new SerializationException("Can't deserialize data [" + Arrays.toString(data) + "] from topic [" + topic + "]", e);
        }
    }
}

코드 조각 7~9 line은 설정하지 않았기 때문에 pass, 실제 type 추론을 10~12 line에서 진행해. JsonDeserializer에서 사용하는 TypeMapper는 DefaultJackson2JavaTypeMapper.class 를 이용하고 있더라고 DefaultJackson2JavaTypeMapper.toJavaType 에서는 header에 클래스 타입을 추론할 수 있는 정보가 있으면 해당 JavaType을 return 하고 없으면 null 을 return 하도록 되어 있어.

하지만 아까 우리는 useHeadersIfPresent를 false 로 설정했기 때문에 header가 비어있어서 JavaType은 null 이 되었고 13~15 line을 통과하고(타입이 있으면 이 라인을 타서 변환돼) 23~25 line을 통해 json을 우리가 처음 생성자에 세팅한 test.orderWorker.kafka.message.CheckedOutMessage 메세지로 정상적으로 변환되는 거였어.

근데 하나 더 궁금한 게 그러면 왜 useHeadersIfPresent = true(기본값) 일 때는 변환이 되지 않았을까? 😇 (제발 멈춰..)

Journey Season 6 (TypeId, who are you?)

아까 거기서 멈췄어야 하는데 여기까지 와버렸어 🤪 이제 돌아갈 수 없어..

useHeadersIfPresent = true(기본값)으로 설정하고 타입 추론 부분에 breakPoint를 걸고 데이터를 확인해 봤어. 아니 근데 header에 내가 넣지도 않은 값이 설정되어 있는 것 아니겠어?

headers.param1 :: "__TypeId__", "test.order.kafka.message.CheckedOutMessage"

??? 이러니까 classNotFoundException이 자꾸 뜨지…

아 누가 넣었는데? 🤬 범인은… 뭐 어차피 너 아니면 나 일 거잖아? 그런데 typeId 값이.. producer 쪽 파일 경로네? 100% producer 가 넣었다. ㄹㅇ 이건 의심의 여지가 없어서 바-로 producer 쪽을 확인해 봤더니 JsonSerializer는 기본 설정으로 header에 __TypeId__를 세팅하더라고?

결국 우리가 사용하던 기본 설정의 JsonSerializer, JsonDeserializer는 타입 매핑을 위해 기본적으로 __TypeId__ 를 헤더에 가지고 다니면서 사용하는 거였지. 그래서 공식 홈페이지에서도 이 값을 이용해서 string pasing으로 추론 타입을 구하는 등의 예시를 작성해둔 거였어.

자 이제 거의 다 된 것 같아! 아까 위에서 놓친 거 하나 빼고 😕

Journey Season Final (ErrorHandlingDeserializer로 무한 루프 막기)

아까 위에서 다시 보자 했던.. Parsing에서 오류 나면 무한 루프 돌면서 에러 나던 거를 막아야겠지…? 기본적으로 Spring에서 Kafka Message를 poll() 한 이후에 Parsing 에러가 나면 하염없이 다시 읽고 다시 실패하고를 반복하게 되어 있어, 이를 방지하고자 Spring에서는 ErrorHandlingDeserializer 사용을 권장해.

Serialization, Deserialization, and Message Conversion :: Spring Kafka

공식 문서에 자세한 설명이 나와 있으니 참고하면 좋아.

우리는 내부적으로 Parsing 오류 시 슬랙으로 메세지를 전송하고 null 처리하고 끝내도록 설계했기 때문에 Code Level에서 ErrorHandlingDeserializer를 만들어 사용하게 했어. 재처리나 Failover 처리를 위해서는 후속 작업이 필요하니 공식 문서를 참고하길 바라.

ErrorHandlingDeserializer<T> errorHandlingDeserializer = new ErrorHandlingDeserializer<>(jsonDeserializer);
errorHandlingDeserializer.setFailedDeserializationFunction(failedDeserializationInfo -> {
    log.error("failed deserialize :: ", failedDeserializationInfo.getException());
    slackSender.sendMessage(failedDeserializationInfo.getException());
    return null;
});

여행을 마무리하며

내 여행의 목표는 크게 3가지였어

  1. objectMapper로 string을 가져다가 변환하는 게 아니라 poll() 하면서 바로 모델로 받고 싶다.
  2. 설정을 최소화 했으면 좋겠다.
  3. 휴먼 에러가 나지 않도록 하고 싶다.

공식 문서에 따르면 Token 방식이나, Key를 지정하는 방식, Header를 직접 파싱 해서 사용하는 방식, Projection 등 여러 가지 방식을 제공하고 있지만 우리 팀에서 사용하기에 가장 적합한 방식은

Header에 있는 TypeId를 무시하고 받은 Message Json을 내부 DTO로 바로 변환하게 하는 설정이었어.

나와 비슷한 고민과 고생을 하는 사람들에게 도움이 되었길 바라. 🥰

references

spring for apache kafka Serialization, Deserialization, and Message Conversion, https://docs.spring.io/spring-kafka/reference/kafka/serdes.html#overview, 2025.01.03