Streams Scala에서 Java API로 마이그레이션하기
Streams Scala에서 Java API로 마이그레이션하기
kafka-streams-scala 라이브러리를 쓰고 있었다면 주의하셔야 할 중요한 소식이 있어요. 이 라이브러리가 Kafka 4.3부터 폐기(deprecated) 되고, Kafka 5.0에서 제거될 예정이에요. 다행히 Java Streams API를 Scala에서 그대로 쓰는 건 크게 어렵지 않아요. 이 페이지에서 어떻게 옮겨가는지 정리해드릴게요.
출처: 문서
본문
⚠️ 폐기 안내:
kafka-streams-scala라이브러리는 Kafka 4.3부터 폐기되었으며 Kafka 5.0에서 제거될 예정이에요. 이 가이드는 Scala 애플리케이션이 Java Streams API를 직접 사용하도록 마이그레이션하는 데 도움을 줄 거예요. 자세한 내용은 KIP-1244를 참고해요.
마이그레이션 개요
Java Streams API는 최소한의 조정만으로 Scala에서 잘 동작해요. 주요 차이점은 다음과 같아요:
- Scala 래퍼 클래스 대신 Java 타입을 직접 사용
StreamsConfig를 통해 Serde를 명시적으로 구성하거나 메서드에 전달
예시: 단어 카운트 애플리케이션
Scala 래퍼 방식 (폐기됨)
import java.util.Properties
import org.apache.kafka.streams.scala.ImplicitConversions._
import org.apache.kafka.streams.scala._
import org.apache.kafka.streams.scala.kstream._
import org.apache.kafka.streams.{KafkaStreams, StreamsConfig}
import org.apache.kafka.streams.scala.serialization.Serdes._
object WordCountScala extends App {
val props = new Properties()
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount")
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
val builder = new StreamsBuilder // Scala wrapper
val textLines: KStream[String, String] = builder.stream[String, String]("input-topic")
val wordCounts: KTable[String, Long] = textLines
.flatMapValues(line => line.toLowerCase.split("\\W+"))
.groupBy((_, word) => word)
.count()
wordCounts.toStream.to("output-topic")
val streams = new KafkaStreams(builder.build(), props)
streams.start()
}
Java API 방식
import java.util.Properties
import org.apache.kafka.streams.{KafkaStreams, StreamsBuilder, StreamsConfig}
import org.apache.kafka.streams.kstream.{KStream, KTable, Produced}
import org.apache.kafka.common.serialization.Serdes
import scala.jdk.CollectionConverters._
object WordCountJava extends App {
val props = new Properties()
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount")
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
// Configure default serdes
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, classOf[Serdes.StringSerde])
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, classOf[Serdes.StringSerde])
val builder = new StreamsBuilder // Java StreamsBuilder
val textLines = builder.stream[String, String]("input-topic")
val wordCounts = textLines
.flatMapValues(_.toLowerCase.split("\\W+"))
.groupBy((_, word) => word)
.count()
wordCounts.toStream.to("output-topic", Produced.`with`(Serdes.String(), Serdes.Long()))
val streams = new KafkaStreams(builder.build(), props)
streams.start()
}
더 알아보기
- Streams 개발자 가이드 — Java API로 스트림즈 앱을 작성하는 전체 문서를 봐요.
- 데이터 타입과 직렬화 — Java API에서 Serde를 구성하는 방법을 봐요.
- KIP-1244 — Scala API 제거 배경을 봐요.