Spark Streaming + Kafka 통합 가이드
Spark Streaming + Kafka 통합 가이드 (Kafka 0.10.0 이상 브로커) (Spark Streaming + Kafka Integration Guide)
Spark Streaming과 Kafka 0.10을 통합하는 방법을 정리한 문서예요. 직접 스트림(Direct Stream)을 만들고, LocationStrategies·ConsumerStrategies, 오프셋 가져오기·저장, SSL/TLS 설정, 배포까지 다뤄요. Kafka 파티션과 Spark 파티션이 1:1로 대응되는 단순한 병렬 처리를 제공해요.
출처: 문서
본문
Kafka 0.10용 Spark Streaming 통합은 단순한 병렬 처리, Kafka 파티션과 Spark 파티션 사이의 1:1 대응, 그리고 오프셋과 메타데이터에 대한 접근을 제공해요. 하지만 새 통합은 simple API 대신 새 Kafka consumer API를 사용하므로, 사용법에 눈에 띄는 차이점이 있어요.
링크 (Linking)
SBT/Maven 프로젝트 정의를 사용하는 Scala/Java 애플리케이션의 경우, 스트리밍 애플리케이션을 다음 아티팩트와 연결하세요 (자세한 내용은 메인 프로그래밍 가이드의 Linking 섹션 참고).
groupId = org.apache.spark
artifactId = spark-streaming-kafka-0-10_2.13
version = 4.2.0
org.apache.kafka 아티팩트(예: kafka-clients)에 의존성을 수동으로 추가하지 마세요. spark-streaming-kafka-0-10 아티팩트가 적절한 전이(transitive) 의존성을 이미 갖고 있으며, 다른 버전은 진단하기 어려운 방식으로 호환되지 않을 수 있어요.
직접 스트림 만들기 (Creating a Direct Stream)
참고로 import의 네임스페이스에 버전이 포함돼요: org.apache.spark.streaming.kafka010
import org.apache.kafka.clients.consumer.ConsumerRecord
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.streaming.kafka010._
import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "localhost:9092,anotherhost:9092",
"key.deserializer" -> classOf[StringDeserializer],
"value.deserializer" -> classOf[StringDeserializer],
"group.id" -> "use_a_separate_group_id_for_each_stream",
"auto.offset.reset" -> "latest",
"enable.auto.commit" -> (false: java.lang.Boolean)
)
val topics = Array("topicA", "topicB")
val stream = KafkaUtils.createDirectStream[String, String](
streamingContext,
PreferConsistent,
Subscribe[String, String](topics, kafkaParams)
)
stream.map(record => (record.key, record.value))
스트림의 각 항목은 ConsumerRecord예요.
import java.util.*;
import org.apache.spark.SparkConf;
import org.apache.spark.TaskContext;
import org.apache.spark.api.java.*;
import org.apache.spark.api.java.function.*;
import org.apache.spark.streaming.api.java.*;
import org.apache.spark.streaming.kafka010.*;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;
import scala.Tuple2;
Map<String, Object> kafkaParams = new HashMap<>();
kafkaParams.put("bootstrap.servers", "localhost:9092,anotherhost:9092");
kafkaParams.put("key.deserializer", StringDeserializer.class);
kafkaParams.put("value.deserializer", StringDeserializer.class);
kafkaParams.put("group.id", "use_a_separate_group_id_for_each_stream");
kafkaParams.put("auto.offset.reset", "latest");
kafkaParams.put("enable.auto.commit", false);
Collection<String> topics = Arrays.asList("topicA", "topicB");
JavaInputDStream<ConsumerRecord<String, String>> stream =
KafkaUtils.createDirectStream(
streamingContext,
LocationStrategies.PreferConsistent(),
ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams)
);
stream.mapToPair(record -> new Tuple2<>(record.key(), record.value()));
가능한 kafkaParams는 Kafka consumer 구성 문서를 참고하세요. Spark 배치 기간이 기본 Kafka 하트비트 세션 타임아웃(30초)보다 길다면 heartbeat.interval.ms와 session.timeout.ms를 적절히 늘리세요. 5분보다 큰 배치의 경우 브로커에서 group.max.session.timeout.ms를 변경해야 해요. 참고로 예제는 enable.auto.commit을 false로 설정하는데, 이에 대한 논의는 아래 Storing Offsets를 참고하세요.
LocationStrategies
새 Kafka consumer API는 메시지를 버퍼에 미리 가져옵니다(pre-fetch). 따라서 성능상 Spark 통합이 executor에 캐시된 consumer를 유지하고(각 배치마다 다시 만드는 대신), 적절한 consumer를 가진 호스트 위치에 파티션을 스케줄링하는 것을 선호하는 것이 중요해요.
대부분의 경우 위에 보여준 것처럼 LocationStrategies.PreferConsistent를 사용해야 해요. 이는 파티션을 사용 가능한 executor들에 고르게 분산해요. executor가 Kafka 브로커와 같은 호스트에 있다면 PreferBrokers를 사용하세요. 이는 해당 파티션의 Kafka 리더에 파티션을 스케줄링하는 것을 선호해요. 마지막으로 파티션 간 부하에 상당한 왜곡(skew)이 있다면 PreferFixed를 사용하세요. 이는 파티션에서 호스트로의 명시적 매핑을 지정할 수 있게 해줘요 (지정되지 않은 파티션은 일관된 위치를 사용해요).
consumer용 캐시는 기본 최대 크기가 64예요. (64 * executor 수)보다 많은 Kafka 파티션을 처리할 것으로 예상된다면 spark.streaming.kafka.consumer.cache.maxCapacity 설정으로 이 값을 변경할 수 있어요.
Kafka consumer 캐싱을 비활성화하려면 spark.streaming.kafka.consumer.cache.enabled를 false로 설정하면 돼요.
캐시는 topicpartition과 group.id로 키가 지정되므로, createDirectStream 호출마다 별도의 group.id를 사용하세요.
ConsumerStrategies
새 Kafka consumer API는 토픽을 지정하는 여러 가지 다른 방법이 있는데, 일부는 객체 인스턴스화 후 상당한 설정이 필요해요. ConsumerStrategies는 Spark가 체크포인트에서 재시작한 후에도 올바르게 구성된 consumer를 얻을 수 있게 해주는 추상화를 제공해요.
위에서 보여준 ConsumerStrategies.Subscribe는 고정된 토픽 컬렉션에 구독할 수 있게 해줘요. SubscribePattern은 regex를 사용해 관심 토픽을 지정할 수 있게 해줘요. 참고로 0.8 통합과 달리 Subscribe나 SubscribePattern을 사용하면 실행 중인 스트림에서 파티션이 추가될 때 응답해야 해요. 마지막으로 Assign은 고정된 파티션 컬렉션을 지정할 수 있게 해줘요. 세 전략 모두 특정 파티션의 시작 오프셋을 지정할 수 있게 해주는 오버로드된 생성자가 있어요.
위 옵션으로 충족되지 않는 특정 consumer 설정 요구가 있다면, ConsumerStrategy는 확장할 수 있는 public 클래스예요.
RDD 만들기 (Creating an RDD)
배치 처리에 더 적합한 사용 사례가 있다면, 정의된 오프셋 범위에 대한 RDD를 만들 수 있어요.
// Import dependencies and create kafka params as in Create Direct Stream above
val offsetRanges = Array(
// topic, partition, inclusive starting offset, exclusive ending offset
OffsetRange("test", 0, 0, 100),
OffsetRange("test", 1, 0, 100)
)
val rdd = KafkaUtils.createRDD[String, String](sparkContext, kafkaParams, offsetRanges, PreferConsistent)
// Import dependencies and create kafka params as in Create Direct Stream above
OffsetRange[] offsetRanges = {
// topic, partition, inclusive starting offset, exclusive ending offset
OffsetRange.create("test", 0, 0, 100),
OffsetRange.create("test", 1, 0, 100)
};
JavaRDD<ConsumerRecord<String, String>> rdd = KafkaUtils.createRDD(
sparkContext,
kafkaParams,
offsetRanges,
LocationStrategies.PreferConsistent()
);
참고로 PreferBrokers는 사용할 수 없어요. 스트림이 없으면 브로커 메타데이터를 자동으로 조회할 드라이버 측 consumer가 없기 때문이에요. 필요하다면 자체 메타데이터 조회로 PreferFixed를 사용하세요.
오프셋 가져오기 (Obtaining Offsets)
stream.foreachRDD { rdd =>
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
rdd.foreachPartition { iter =>
val o: OffsetRange = offsetRanges(TaskContext.get.partitionId)
println(s"${o.topic} ${o.partition} ${o.fromOffset} ${o.untilOffset}")
}
}
stream.foreachRDD(rdd -> {
OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
rdd.foreachPartition(consumerRecords -> {
OffsetRange o = offsetRanges[TaskContext.get().partitionId()];
System.out.println(
o.topic() + " " + o.partition() + " " + o.fromOffset() + " " + o.untilOffset());
});
});
참고로 HasOffsetRanges로의 타입 캐스트는 createDirectStream의 결과에 호출된 첫 번째 메서드에서만 성공하며, 메서드 체인의 나중에서 는 안 돼요. RDD 파티션과 Kafka 파티션 사이의 1:1 매핑은 reduceByKey()나 window()처럼 shuffle이나 repartition을 하는 메서드 후에는 유지되지 않는다는 점을 알아두세요.
오프셋 저장 (Storing Offsets)
실패 시의 Kafka 전달 의미론(delivery semantics)은 오프셋을 어떻게, 언제 저장하느냐에 달려 있어요. Spark 출력 연산은 최소 한 번(at-least-once)이에요. 따라서 정확히 한 번(exactly-once) 의미론과 동등한 것을 원한다면, 멱등(idempotent) 출력 후에 오프셋을 저장하거나, 출력과 함께 원자적 트랜잭션으로 오프셋을 저장해야 해요. 이 통합에서는 오프셋을 저장하는 방법으로 신뢰성(및 코드 복잡성)이 증가하는 순서로 3가지 옵션이 있어요.
체크포인트 (Checkpoints)
Spark 체크포인팅을 활성화하면 오프셋이 체크포인트에 저장돼요. 활성화하기 쉽지만 단점이 있어요. 반복된 출력이 발생하므로 출력 연산이 멱등(idempotent)이어야 해요. 트랜잭션은 옵션이 아니에요. 또한 애플리케이션 코드가 변경되면 체크포인트에서 복구할 수 없어요. 계획된 업그레이드의 경우 새 코드를 이전 코드와 동시에 실행해 완화할 수 있어요 (어차피 출력은 멱등이어야 하므로 충돌하지 않아야 해요). 하지만 코드 변경이 필요한 계획되지 않은 실패의 경우, 알려진 양호한 시작 오프셋을 식별할 다른 방법이 없다면 데이터를 잃게 돼요.
Kafka 자체 (Kafka itself)
Kafka는 특수한 Kafka 토픽에 오프셋을 저장하는 오프셋 커밋(offset commit) API가 있어요. 기본적으로 새 consumer는 주기적으로 오프셋을 자동 커밋해요. 이는 거의 확실히 여러분이 원하는 것이 아닌데, consumer가 성공적으로 폴링한 메시지가 아직 Spark 출력 연산으로 이어지지 않았을 수 있어 정의되지 않은 의미론을 초래하기 때문이에요. 그래서 위의 스트림 예제가 "enable.auto.commit"을 false로 설정한 것이에요. 하지만 출력이 저장되었음을 확인한 후 commitAsync API를 사용해 Kafka에 오프셋을 커밋할 수 있어요. 체크포인트에 비해 이점은 Kafka가 애플리케이션 코드의 변경과 무관하게 내구성 있는 저장소라는 점이에요. 하지만 Kafka는 트랜잭션이 아니므로 출력은 여전히 멱등이어야 해요.
stream.foreachRDD { rdd =>
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
// some time later, after outputs have completed
stream.asInstanceOf[CanCommitOffsets].commitAsync(offsetRanges)
}
HasOffsetRanges와 마찬가지로 CanCommitOffsets로의 캐스트는 변환 후가 아니라 createDirectStream의 결과에 호출된 경우에만 성공해요. commitAsync 호출은 스레드 안전(threadsafe)하지만, 의미 있는 의미론을 원한다면 출력 후에 발생해야 해요.
stream.foreachRDD(rdd -> {
OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
// some time later, after outputs have completed
((CanCommitOffsets) stream.inputDStream()).commitAsync(offsetRanges);
});
자신만의 데이터 저장소 (Your own data store)
트랜잭션을 지원하는 데이터 저장소의 경우, 결과와 같은 트랜잭션에서 오프셋을 저장하면 실패 상황에서도 둘을 동기화할 수 있어요. 반복되거나 건너뛴 오프셋 범위를 감지하는 데 주의한다면, 트랜잭션을 롤백해 중복·누락 메시지가 결과에 영향을 주는 것을 방지할 수 있어요. 이는 정확히 한 번(exactly-once) 의미론과 동등한 것을 제공해요. 보통 멱등으로 만들기 어려운 집계 결과에서 비롯된 출력에도 이 전술을 사용할 수 있어요.
// The details depend on your data store, but the general idea looks like this
// begin from the offsets committed to the database
val fromOffsets = selectOffsetsFromYourDatabase.map { resultSet =>
new TopicPartition(resultSet.string("topic"), resultSet.int("partition")) -> resultSet.long("offset")
}.toMap
val stream = KafkaUtils.createDirectStream[String, String](
streamingContext,
PreferConsistent,
Assign[String, String](fromOffsets.keys.toList, kafkaParams, fromOffsets)
)
stream.foreachRDD { rdd =>
val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges
val results = yourCalculation(rdd)
// begin your transaction
// update results
// update offsets where the end of existing offsets matches the beginning of this batch of offsets
// assert that offsets were updated correctly
// end your transaction
}
// The details depend on your data store, but the general idea looks like this
// begin from the offsets committed to the database
Map<TopicPartition, Long> fromOffsets = new HashMap<>();
for (resultSet : selectOffsetsFromYourDatabase)
fromOffsets.put(new TopicPartition(resultSet.string("topic"), resultSet.int("partition")), resultSet.long("offset"));
}
JavaInputDStream<ConsumerRecord<String, String>> stream = KafkaUtils.createDirectStream(
streamingContext,
LocationStrategies.PreferConsistent(),
ConsumerStrategies.<String, String>Assign(fromOffsets.keySet(), kafkaParams, fromOffsets)
);
stream.foreachRDD(rdd -> {
OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
Object results = yourCalculation(rdd);
// begin your transaction
// update results
// update offsets where the end of existing offsets matches the beginning of this batch of offsets
// assert that offsets were updated correctly
// end your transaction
});
SSL / TLS
새 Kafka consumer는 SSL을 지원해요. 활성화하려면 createDirectStream / createRDD에 전달하기 전에 kafkaParams를 적절히 설정하세요. 참고로 이는 Spark와 Kafka 브로커 사이의 통신에만 적용돼요. Spark 노드 간 통신을 별도로 보호하는 것은 여전히 여러분의 책임이에요.
val kafkaParams = Map[String, Object](
// the usual params, make sure to change the port in bootstrap.servers if 9092 is not TLS
"security.protocol" -> "SSL",
"ssl.truststore.location" -> "/some-directory/kafka.client.truststore.jks",
"ssl.truststore.password" -> "test1234",
"ssl.keystore.location" -> "/some-directory/kafka.client.keystore.jks",
"ssl.keystore.password" -> "test1234",
"ssl.key.password" -> "test1234"
)
Map<String, Object> kafkaParams = new HashMap<String, Object>();
// the usual params, make sure to change the port in bootstrap.servers if 9092 is not TLS
kafkaParams.put("security.protocol", "SSL");
kafkaParams.put("ssl.truststore.location", "/some-directory/kafka.client.truststore.jks");
kafkaParams.put("ssl.truststore.password", "test1234");
kafkaParams.put("ssl.keystore.location", "/some-directory/kafka.client.keystore.jks");
kafkaParams.put("ssl.keystore.password", "test1234");
kafkaParams.put("ssl.key.password", "test1234");
배포 (Deploying)
다른 Spark 애플리케이션과 마찬가지로 spark-submit을 사용해 애플리케이션을 실행해요.
Scala/Java 애플리케이션의 경우 SBT나 Maven을 프로젝트 관리에 사용한다면 spark-streaming-kafka-0-10_2.13과 그 의존성을 애플리케이션 JAR에 패키징하세요. spark-core_2.13과 spark-streaming_2.13은 Spark 설치에 이미 있으므로 provided 의존성으로 표시하세요. 그런 다음 spark-submit으로 애플리케이션을 실행하세요 (메인 프로그래밍 가이드의 Deploying 섹션 참고).
보안 (Security)
Structured Streaming Security를 참고하세요.
추가 주의사항 (Additional Caveats)
- Kafka 네이티브 싱크(sink)는 사용할 수 없으므로 위임 토큰(delegation token)은 consumer 측에서만 사용돼요.
더 알아보기 (Learn more)
- 아파치 스파크 Kafka 통합 가이드 (원문)
- Structured Streaming + Kafka 통합 가이드 — Structured Streaming용 Kafka 통합
- Spark Streaming 프로그래밍 가이드 (원문) — Spark Streaming 가이드