데이터 타입과 직렬화
데이터 타입과 직렬화
카프카는 데이터를 바이트로 주고받기 때문에, Kafka Streams 애플리케이션은 레코드 키와 값의 데이터 타입에 대한 Serde(Serializer/Deserializer)를 반드시 제공해야 해요. 이 페이지에서는 기본 Serde를 설정하는 방법, 특정 지점에서 오버라이드하는 방법, 제공되는 내장 Serde들, 그리고 커스텀 Serde를 만드는 방법까지 정리해드릴게요.
출처: 문서
본문
모든 Kafka Streams 애플리케이션은 필요한 시점에 데이터를 구체화(materialize)하기 위해 레코드 키와 레코드 값의 데이터 타입(예: java.lang.String)에 대한 Serde를 제공해야 해요. 그러한 Serde 정보가 필요한 연산에는 stream(), table(), to(), repartition(), groupByKey(), groupBy()가 있어요.
다음 두 가지 방법 중 하나로 Serde를 제공할 수 있으며, 최소 하나는 사용해야 해요:
java.util.Properties구성 인스턴스에 기본 Serde를 설정.- 적절한 API 메서드를 호출할 때 명시적 Serde를 지정해 기본값을 오버라이드.
Serde 구성하기
Streams 구성에 지정된 Serde는 Kafka Streams 애플리케이션에서 기본값으로 사용돼요. 이 구성의 기본값이 null이므로, 이 구성을 사용해 기본 Serde를 설정하거나 아래에 설명된 대로 Serde를 명시적으로 전달해야 해요.
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.StreamsConfig;
Properties settings = new Properties();
// Default serde for keys of data records (here: built-in serde for String type)
settings.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
// Default serde for values of data records (here: built-in serde for Long type)
settings.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass().getName());
기본 Serde 오버라이드
적절한 API 메서드에 Serde를 전달해 명시적으로 지정할 수도 있으며, 이는 기본 serde 설정을 오버라이드해요:
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
final Serde<String> stringSerde = Serdes.String();
final Serde<Long> longSerde = Serdes.Long();
// The stream userCountByRegion has type `String` for record keys (for region)
// and type `Long` for record values (for user counts).
KStream<String, Long> userCountByRegion = ...;
userCountByRegion.to("RegionCountsTopic", Produced.with(stringSerde, longSerde));
일부 필드에 대해 기본값을 유지하면서 serde를 선택적으로 오버라이드하고 싶다면, 기본 설정을 활용하고 싶을 때 serde를 지정하지 않으면 돼요:
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
// Use the default serializer for record keys (here: region as String) by not specifying the key serde,
// but override the default serializer for record values (here: userCount as Long).
final Serde<Long> longSerde = Serdes.Long();
KStream<String, Long> userCountByRegion = ...;
userCountByRegion.to("RegionCountsTopic", Produced.valueSerde(Serdes.Long()));
들어오는 레코드 중 일부가 손상되었거나 형식이 잘못되었다면, 디시리얼라이저 클래스가 오류를 보고하게 될 거예요. 1.0.x부터 DeserializationExceptionHandler 인터페이스를 도입해 그러한 레코드를 어떻게 처리할지 커스터마이즈할 수 있어요. 인터페이스의 커스터마이즈 구현은 StreamsConfig를 통해 지정할 수 있어요. 자세한 내용은 Streams 애플리케이션 구성하기 섹션을 참고해요.
사용 가능한 Serdes
기본 타입과 기본적인 타입
Apache Kafka는 kafka-clients Maven 아티팩트에 Java 기본 타입과 byte[] 같은 기본 타입을 위한 몇 가지 내장 serde 구현을 포함해요:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>4.3.1</version>
</dependency>
이 아티팩트는 org.apache.kafka.common.serialization 패키지 아래 다음 serde 구현을 제공하며, 예를 들어 Streams 구성에서 기본 시리얼라이저를 정의할 때 활용할 수 있어요.
| 데이터 타입 | Serde |
|---|---|
byte[] |
Serdes.ByteArray(), Serdes.Bytes() (아래 팁 참조) |
| ByteBuffer | Serdes.ByteBuffer() |
| Double | Serdes.Double() |
| Integer | Serdes.Integer() |
| Long | Serdes.Long() |
| String | Serdes.String() |
| UUID | Serdes.UUID() |
| Void | Serdes.Void() |
| List | Serdes.ListSerde() |
| Boolean | Serdes.Boolean() |
팁:
Bytes는 올바른 동등성(equality)과 순서(ordering) 의미론을 지원하는 Java의byte[](바이트 배열) 래퍼예요. 애플리케이션에서byte[]대신Bytes사용을 고려해볼 만해요.
JSON
Kafka Streams 코드 예제에는 JSON을 위한 기본 serde 구현도 포함돼요: PageViewTypedDemo. 예시에 나와 있듯이, JSONSerdes 내부 클래스 Serdes.serdeFrom(<serializerInstance>, <deserializerInstance>)를 사용해 JSON 호환 시리얼라이저·디시리얼라이저를 구성할 수 있어요.
윈도우 Serdes
Apache Kafka Streams는 kafka-streams Maven 아티팩트에 윈도우 타입을 위한 serde 구현을 포함해요:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>4.3.1</version>
</dependency>
이 아티팩트는 org.apache.kafka.streams.kstream 패키지 아래 다음 윈도우 serde 구현을 제공해요:
Serdes:
WindowedSerdes.TimeWindowedSerde<T>WindowedSerdes.SessionWindowedSerde<T>
Serializers:
TimeWindowedSerializer<T>SessionWindowedSerializer<T>
Deserializers:
TimeWindowedDeserializer<T>SessionWindowedDeserializer<T>
코드에서의 사용
애플리케이션 코드에서 윈도우 serde를 사용할 때는 보통 생성자나 팩토리 메서드를 통해 인스턴스를 만들어요:
// Time windowed serde - using factory method
Serde<Windowed<String>> timeWindowedSerde =
WindowedSerdes.timeWindowedSerdeFrom(String.class, 500L);
// Time windowed serde - using constructor
Serde<Windowed<String>> timeWindowedSerde2 =
new WindowedSerdes.TimeWindowedSerde<>(Serdes.String(), 500L);
// Session windowed serde - using factory method
Serde<Windowed<String>> sessionWindowedSerde =
WindowedSerdes.sessionWindowedSerdeFrom(String.class);
// Session windowed serde - using constructor
Serde<Windowed<String>> sessionWindowedSerde2 =
new WindowedSerdes.SessionWindowedSerde<>(Serdes.String());
// Using individual serializers/deserializers
TimeWindowedSerializer<String> serializer = new TimeWindowedSerializer<>(Serdes.String().serializer());
TimeWindowedDeserializer<String> deserializer = new TimeWindowedDeserializer<>(Serdes.String().deserializer(), 500L);
명령줄에서의 사용
bin/kafka-console-consumer.sh 같은 명령줄 도구를 사용할 때는 구성 속성을 통해 내부 클래스와 윈도우 크기를 전달해 윈도우 디시리얼라이저를 구성할 수 있어요. 속성 이름은 접두사 패턴을 사용해요:
# Time windowed deserializer configuration
--formatter-property print.key=true \
--formatter-property key.deserializer=org.apache.kafka.streams.kstream.TimeWindowedDeserializer \
--formatter-property key.deserializer.windowed.inner.deserializer.class=org.apache.kafka.common.serialization.StringDeserializer \
--formatter-property key.deserializer.window.size.ms=500
# Session windowed deserializer configuration
--formatter-property print.key=true \
--formatter-property key.deserializer=org.apache.kafka.streams.kstream.SessionWindowedDeserializer \
--formatter-property key.deserializer.windowed.inner.deserializer.class=org.apache.kafka.common.serialization.StringDeserializer
폐기된 구성
다음 StreamsConfig 파라미터는 시리얼라이저/디시리얼라이저 생성자에 파라미터를 직접 전달하는 방식으로 대체되어 폐기되었어요:
StreamsConfig.WINDOWED_INNER_CLASS_SERDE는TimeWindowedSerializer.WINDOWED_INNER_SERIALIZER_CLASS와TimeWindowedDeserializer.WINDOWED_INNER_DESERIALIZER_CLASS를 위해 폐기되었어요.StreamsConfig.WINDOW_SIZE_MS_CONFIG는TimeWindowedDeserializer.WINDOW_SIZE_MS_CONFIG를 위해 폐기되었어요.
커스텀 Serdes 구현하기
커스텀 Serdes를 구현해야 한다면, 기존 Serdes의 소스 코드 참조(이전 섹션)를 살펴보는 것이 가장 좋은 시작점이에요. 일반적으로 워크플로는 다음과 비슷할 거예요:
org.apache.kafka.common.serialization.Serializer를 구현해 데이터 타입T에 대한 시리얼라이저를 작성.org.apache.kafka.common.serialization.Deserializer를 구현해T에 대한 디시리얼라이저를 작성.org.apache.kafka.common.serialization.Serde를 구현해T에 대한 serde를 작성. 수동으로 하거나(이전 섹션의 기존 Serdes 참조)Serdes.serdeFrom(Serializer<T>, Deserializer<T>)같은 Serdes의 헬퍼 함수를 활용.
참고로, KafkaStreams에 제공되는 구성에서 커스텀 serde를 사용하려면 (제네릭 타입이 없는) 자신의 클래스를 구현해야 해요. serde 클래스에 제네릭 타입이 있거나 Serdes.serdeFrom(Serializer<T>, Deserializer<T>)를 사용한다면, 메서드 호출을 통해서만 serde를 전달할 수 있어요 (예: builder.stream("topicName", Consumed.with(...))).
Kafka Streams DSL for Scala 암묵적(Implicit) Serdes
Kafka Streams DSL for Scala를 사용할 때는 기본 Serdes를 구성할 필요가 없어요. 사실 구성하는 것이 지원되지 않아요. 대신 Serdes는 일반적인 기본 데이터 타입에 대한 기본 구현으로 암묵적으로 제공돼요. 자세한 내용은 DSL API 문서의 Implicit Serdes와 User-Defined Serdes 섹션을 참고해요.
더 알아보기
- Streams 애플리케이션 구성하기 — Serde와 관련된 Streams 구성 전체를 봐요.
- Streams DSL — DSL에서 Serde를 사용하는 패턴을 봐요.
- Streams 애플리케이션 테스트하기 — Serde가 잘 동작하는지 테스트해요.