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()
}

더 알아보기