Data
CDC 파이프라인 정합성 검사 Spark 잡 개발 - Part 1. 코드 설계편
dawn.choi카카오
2025년 7월 28일
원문에서 보기 ↗1. 들어가며
안녕하세요, 카카오 데이터분석플랫폼 조직의 dawn(최여명)입니다.
저희 조직에서는 카카오 내 여러 서비스의 데이터를 한 곳으로 응집하기 위해 약 300개 이상의 대규모 CDC(Change Data Capture) 파이프라인을 운영하고 있습니다. 작년 사내 기술 컨퍼런스인 이프카카오에서도 대용량 데이터베이스 동기화를 위한 최적의 CDC 시스템 구축기라는 주제로 발표를 진행한 적이 있는데요. 이번 글에서는 발표에서 간략히 언급했던 CDC 파이프라인 데이터 정합성 검사에 대해 보다 상세히 소개하고자 합니다.
정합성 검사는 사내에 제공하는 데이터의 신뢰성을 보장하기 위해 매우 중요한데요. 운영 중인 파이프라인에서는 데이터가 유실되거나, NULL 여부, 타입 차이로 인한 불일치 등 다양한 이슈가 발생할 가능성이 항상 존재하기 때문입니다. 특히 신규 연동 파이프라인의 경우, 연동 직후 데이터가 올바르게 적재되었는지 철저하게 검증하는 것이 필수적입니다. 이렇게 중요한 정합성 검사를 위한 Spark 잡 코드를 어떻게 설계하고 성능 최적화를 진행했는지에 대한 경험을 공유드리고자 합니다.
이번 글은 두 개의 파트로 나누어 작성할 예정입니다.
- Part 1. 코드 설계편에서는 Scala 기반 Spark 잡 코드를 바닥부터 어떻게 설계했는지를 중심으로 다룹니다. 구체적으로 컴포넌트를 어떻게 상속하고 분리하여 유연하고 확장 가능한 구조를 만들었는지, 또 코드 품질과 유지보수성을 높이기 위해 어떤 디자인 패턴을 적용했는지에 대해 설명합니다.
- Part 2. Spark 최적화편에서는 Spark가 데이터베이스에서 데이터를 읽어올 때 성능을 높이기 위해 적용한 다양한 최적화 기법과 팁들을 자세히 소개할 예정입니다.
그럼 먼저 Part 1인 코드 설계편부터 시작해보겠습니다.
2. 잡 제출 JSON 정의
정합성 검사를 위한 잡(Job)은 사전에 정의된 JSON 형태로 제출합니다. 이 JSON에는 소스 테이블과 타겟 테이블의 접속 정보, 스캔 모드, 리포트 방식, Spark 설정 등 다양한 설정 값이 포함됩니다. 각 JSON은 테이블 단위로 구성되는데요. 이렇게 최소 단위로 나누는 이유는 이슈가 발생했을 때 개별 테이블 단위로 감지할 수 있고, 하나의 잡이 실패하더라도 다른 잡에 영향을 주지 않도록 하기 위함입니다.
쉬운 이해를 위해 제출하는 JSON 예를 보여드리면 아래와 같습니다.
{
"source": {
"datasourceType": "mysql",
"host": "hostname-slv.mydb.daumkakao.io",
"database": "database_name",
"table": "table_name",
"username": "username",
"password": "VAULT(...)",
"pkColumns": ["id"],
"numPartitions": 10
},
"target": {
"datasourceType": "iceberg",
"catalog": "bi_iceberg_catalog",
"namespace": "database_name",
"table": "table_name",
"pkColumns": ["id"]
},
"report": {
"reportTypes": ["database"]
},
"mode": {
"scanType": "fullscan",
"sparkConfig": {
"spark.executor.memory": "4g",
"spark.executor.instances": "8",
"spark.executor.cores": "4"
}
}
}
이 JSON 구조에 맞춰, 각 필드를 Scala case class로 정의해두고 있습니다. 이렇게 하면 설정 항목의 타입과 포맷을 명확하게 관리할 수 있고, 어떤 옵션들이 있는지도 쉽게 파악할 수 있어 개발자 입장에서 매우 유용합니다.
case class Config(
source: DatasourceConfig,
target: DatasourceConfig,
report: ReportConfig,
mode: ModeConfig
)
특히 DatasourceConfig는 공통 필드를 가지는 추상 타입(trait)으로 정의되어 있고, 실제 구현체는 MySQL, MongoDB, Iceberg 등 데이터소스 종류에 따라 다르게 구성됩니다. 또한 JSON을 해당 구현체로 매핑하기 위해 Jackson의 다형성 직렬화 기능을 활용했습니다.
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "datasourceType")
@JsonSubTypes(
Array(
new JsonSubTypes.Type(value = classOf[MysqlDatasourceConfig], name = "mysql"),
new JsonSubTypes.Type(value = classOf[MongoDatasourceConfig], name = "mongodb"),
new JsonSubTypes.Type(value = classOf[IcebergDatasourceConfig], name = "iceberg")
)
)
trait DatasourceConfig {
def datasourceType: DatasourceType
def pkColumns: Array[String]
def database: String
def table: String
def accessPoint: String
def keyColumn: String
def excludedColumns: Set[String]
def translators: Array[String]
}
case class MysqlDatasourceConfig(
datasourceType: DatasourceType = DatasourceType.Mysql,
host: String,
port: Int = 3306,
database: String,
table: String,
pkColumns: Array[String],
username: String,
password: String,
customSchema: String = "",
numPartitions: Int = 1,
keyColumn: String = "",
excludedColumns: Set[String] = Set.empty,
translators: Array[String] = Array.empty[String],
jdbcParameters: Map[String, String] = Map.empty
) extends DatasourceConfig {
override def accessPoint: String = host
}
// MongoDatasourceConfig와 IcebergDatasourceConfig는 생략
위와 같이 @JsonTypeInfo와 @JsonSubTypes 어노테이션을 사용하면 datasourceType 필드의 값에 따라 적절한 자식 클래스가 자동으로 매핑되어 역직렬화됩니다. 이를 통해 JSON을 파싱할 때 데이터소스 종류에 따라 다른 설정 객체를 쉽게 생성할 수 있습니다. 이런 방식으로 필수 필드는 공통 인터페이스에 정의하고, 각 데이터소스마다 필요한 필드를 개별 case class에 선언하여 사용하고 있습니다. 이를 통해 JSON 구조와 코드 간의 일관성을 유지하고, 추후 확장도 용이하게 만들 수 있습니다.
3. 컴포넌트 분리 및 설계
정합성 검사가 수행해야 할 역할은 크게 분리하자면 네가지입니다.
- 소스 테이블과 타겟 테이블에서 데이터를 읽고 (Reader),
- DB 기종의 차이나 보안 데이터 유무에 따른 값 차이를 보정해주고 (Translater),
- 두 데이터셋 간에 차이가 있는지 비교하고 (Comparator),
- 비교 결과를 어딘가에 보고한다. (Reporter)
이 네 단계에 맞춰 코드 구조를 모듈화하고, 각 단계별로 독립적인 컴포넌트로 분리했습니다. 각 컴포넌트는 자신의 역할에만 집중하면서, 다음 단계로 필요한 정보만을 깔끔하게 전달합니다. 또한 코드의 일관성과 확장성을 위해, 상위 추상 클래스에서 공통 동작을 정의하고, 하위 클래스에서는 케이스에 따라 필요한 동작을 오버라이드하도록 구현했습니다. 이 구조를 통해 데이터소스나 리포트 방식이 추가되거나 로직이 변경되는 경우에도 쉽게 개선하고 확장 가능합니다.
Reader
Reader는 소스 테이블과 타겟 테이블에서 데이터를 읽는 역할을 수행하는 컴포넌트입니다. 이때 중요한 점은 어떤 종류의 데이터베이스에서 데이터를 읽는가입니다. 예를 들어, MySQL 테이블과 Iceberg 테이블은 Spark에서 사용하는 포맷과 옵션이 다르기 때문에, 데이터소스에 따라 서로 다른 처리 방식이 필요합니다. 또한 데이터를 읽는 방식에도 여러 가지 모드가 존재하는데요.
- 전체 데이터를 읽는
fullscan - 특정 PK(Primary Key)만 조회하는
keybased - 한정된 범위만 조회하는
limit
이러한 모드는 동작 방식이 데이터소스에 따라 다르며, 이를 공통화하면서도 유연하게 구현하기 위해 DatasourceReader라는 추상 클래스를 정의했습니다. 그리고 각 데이터소스에 맞는 구현체로 MysqlReader, IcebergReader등을 생성하여 세부 내용을 구현했습니다.
abstract class DatasourceReader(datasourceConfig: DatasourceConfig, modeConfig: ModeConfig) {
def readWithFullScanMode(): DataFrame
def readWithKeyBasedMode(): DataFrame
def readWithLimitMode(): DataFrame
}
class MysqlReader(datasourceConfig: MysqlDatasourceConfig, modeConfig: ModeConfig) extends DatasourceReader(datasourceConfig, modeConfig) {
override def readWithFullScanMode(): DataFrame = {
// MySQL에서 사용되는 fullscan 모드 로직 구현
}
override def readWithKeyBasedMode(): DataFrame = {
// MySQL에서 사용되는 keybased 모드 로직 구현
}
override def readWithLimitMode(): DataFrame = {
// MySQL에서 사용되는 limit 모드 로직 구현
}
}
class IcebergReader(datasourceConfig: IcebergDatasourceConfig, modeConfig: ModeConfig) extends DatasourceReader(datasourceConfig, modeConfig) {
override def readWithFullScanMode(): DataFrame = {
// Iceberg에서 사용되는 fullscan 모드 로직 구현
}
override def readWithKeyBasedMode(): DataFrame = {
// Iceberg에서 사용되는 keybased 모드 로직 구현
}
override def readWithLimitMode(): DataFrame = {
// Iceberg에서 사용되는 limit 모드 로직 구현
}
}
이러한 구조를 통해 데이터소스별로 필수적인 스캔 로직 구현을 강제할 수 있고, 새로운 데이터소스나 스캔 모드가 추가되더라도 최소한의 코드 변경만으로 쉽게 확장할 수 있습니다.
Translator
Translator는 데이터를 읽은 후, 두 데이터셋을 비교하기 전에 값 변환이 필요할 때 사용하는 컴포넌트입니다. CDC 연동시 데이터소스가 서로 다를 경우 데이터 타입의 차이로 인해 특정 컬럼의 실제 값이 달라질 수 있는데요. 또한, 보안존 내 DB와 연동할 경우 특정 컬럼이 해싱 또는 마스킹 처리되어 값이 변경될 수도 있습니다. 이러한 차이는 불가피하기 때문에 비교 전에 값 변환 작업이 반드시 필요합니다.
예를 들어, MySQL에서 Iceberg로 데이터를 적재할 때 Iceberg는 JSON 타입을 지원하지 않아 문자열 형태로 데이터를 저장하게 됩니다. 이때 JSON 데이터의 포맷이나 공백 등 세부적인 표현이 달라질 수 있는데요. 따라서 비교를 수행하기 전에 양쪽 데이터셋의 JSON 포맷을 표준화된 형태의 문자열로 변환한 뒤 비교합니다.
Translator는 모두 발동 조건과 로직이 다르기 때문에, Translator라는 추상 클래스를 정의하고 구체적인 로직은 하위 클래스에서 구현하도록 했습니다. 또한 데이터소스의 특성상 기본적으로 적용해야 하는 Translator가 있고, 해싱이나 마스킹과 같이 사용자가 설정에 따라 추가하는 Translator가 존재하기 때문에 두 가지 Translator를 병합하여 리스트로 관리하고 데이터에 차례대로 적용하여 변환을 수행합니다.
아래는 Translator를 로딩하고 적용하는 예시 코드입니다.
def loadTranslators(
sourceDfWrapper: DataFrameWrapper,
targetDfWrapper: DataFrameWrapper,
config: Config
): Unit = {
// datasource 타입에 따른 default translator 설정
val defaultSourceTranslators =
TranslatorFactory.getDefaultTranslatorStringList(
config.source.datasourceType,
config.target.datasourceType,
TranslatorType.Source
)
val defaultTargetTranslators =
TranslatorFactory.getDefaultTranslatorStringList(
config.source.datasourceType,
config.target.datasourceType,
TranslatorType.Target
)
// default translator와 사용자가 json config에 제출한 custom translator 병합
val sourceTranslators = defaultSourceTranslators ++ config.source.translators
val targetTranslators = defaultTargetTranslators ++ config.target.translators
// translate에 사용할 부가정보 로딩
val sourceTranslatorMap = TranslatorFactory.loadDefaultTranslatorMap()
val targetTranslatorMap = TranslatorFactory.loadDefaultTranslatorMap()
// 병합된 default + custom translator 수행
sourceTranslators.foreach(translator => {
val translatorInstance = TranslatorFactory.getTranslator(translator)
translatorInstance.translate(sourceDFWrapper, sourceTranslatorMap)
})
// targetTranslators도 이어서 수행...
}
예시로 추상 클래스 Translator를 상속해 JSON 컬럼을 변환하는 역할을 수행하는 JsonTranslator는 다음과 같이 구현됩니다.
class JsonTranslator extends Translator {
override def translate(dfWrapper: DataFrameWrapper, extraParams: Map[String, Any] = Map.empty): Unit = {
logger.info("Translating JSON")
val jsonColumns = extraParams.get("columns") match {
case Some(columns: Seq[String]) => columns
case _ => Seq.empty[String]
}
val resultDF = convertJsonType(jsonColumns, dfWrapper.getDF)
dfWrapper.setDF(resultDF)
}
}
Comparator
Comparator는 두 데이터셋을 비교하는 컴포넌트입니다. 소스와 타겟 데이터셋에서 동일한 기본 키(PK, Primary Key)를 가진 레코드 간의 값을 비교 하는 역할을 수행합니다. 비교할 데이터셋은 앞선 컴포넌트에서 Spark의 DataFrame 형태로 추상화되고 필요한 값 변환도 완료된 상태이기 때문에 비교 로직은 단일 구현으로 충분합니다. 따라서 해당 컴포넌트는 DataComparator라는 하나의 클래스로 구성했습니다.
DataComparator는 소스 데이터와 타겟 데이터, 그리고 비교에 사용할 기본 키 컬럼 목록을 입력으로 받아, 비교 결과를 객체 내부에 저장하는 명령 기반의 상태 유지형 컴포넌트 입니다. 비교는 한 번의 명령형 단일 메소드 호출로 수행되며 호출 시 결과가 내부 상태로 저장되고 이후 별도의 메서드로 결과를 참조할 수 있도록 설계했습니다. 이 구조는 비교 연산을 하나의 컨텍스트로 묶고, 실행 시점과 결과 참조 시점을 명확하게 분리하는 데 용이합니다.
아래는 예시 코드입니다.
class DataComparator(source: DataFrame, target: DataFrame, pkColumns: Array[String]) {
private var resultDF: DataFrame = _
private var totalCount: Long = _
def compare(): Unit = {
// 비교를 위해 source, target을 full outer join
val joinedDF = source
.alias("source")
.join(target.alias("target"), pkColumns, "full_outer")
totalCount = joinedDF.count()
// 비교 로직 구현... (생략)
// resultDF에는 최종적으로 일치하지 않은 레코드만 담음
val filteredDF = joinedDF.filter(pkNullCondition || nonPkMismatchCondition)
resultDF = filteredDF
}
def getTotalCount: Long = totalCount
def getMatchedCount: Long = totalCount - resultDF.count()
def getMismatchedCount: Long = resultDF.count()
}
Reporter
Reporter는 비교 결과를 원하는 채널로 전달하는 컴포넌트입니다. 콘솔 출력, 사내 메신저 전송, 데이터베이스 저장 등 다양한 방식으로 결과를 보고할 수 있으며 사용자는 잡 제출 시 어떤 채널로 리포트를 전달할지 선택할 수 있습니다.
채널마다 리포트 방식이 다르고 새로운 채널이 필요할시 유연하게 추가하기 위해, Reporter라는 추상 클래스를 정의하고, ConsoleReporter, DatabaseReporter 등 각 채널별 리포트 로직은 하위 클래스에서 구현하도록 설계했습니다. 또한 비교 결과는 ConsistencyResult라는 DTO(Data Transfer Object) 객체로 전달되며 결과 포맷의 통일성을 위해 getFormattedResult 함수는 부모 클래스인 Reporter에서 공통으로 구현되어 모든 하위 클래스에서 재사용됩니다.
case class ConsistencyResult(
startTime: Long,
totalCount: Long,
mismatchedCount: Long,
config: Config
)
abstract class Reporter(result: ConsistencyResult) {
protected val accuracy: Double = getAccuracyPercentage()
def report(): Unit
protected def getFormattedResult: String = {
// 결과 스트링 예쁘게 포맷팅하여 반환
}
private def getAccuracyPercentage(scale: Int = 2): Double = {
// 변경분이 존재하지 않는 경우 정확도는 100%로 처리
if (result.totalCount == 0) {
return 100
}
if (result.totalCount < result.mismatchedCount) {
throw new IllegalArgumentException("Total count cannot be less than mismatched count")
}
val accuracy = (result.totalCount - result.mismatchedCount).toDouble / result.totalCount * 100
val formattedAccuracy = BigDecimal(accuracy).setScale(scale, BigDecimal.RoundingMode.HALF_UP).toDouble
formattedAccuracy
}
4. 적용한 디자인 패턴
디자인 패턴은 소프트웨어를 설계할 때 자주 반복되는 문제들에 대한 해결책을 미리 정리해 둔 방법입니다. 적절한 디자인 패턴 활용은 개발자들이 개별 클래스나 객체 수준보다 더 추상화된 수준에서 코드 구조를 이해할 수 있도록 돕습니다. 이를 통해 코드의 구조가 더 견고해지고, 개발자 간 커뮤니케이션 또한 더 명확해집니다.
이 프로젝트에서도 다양한 데이터소스와 실행 환경을 고려하여 여러 디자인 패턴을 적극적으로 도입했는데 적용한 주요 디자인 패턴을 하나씩 소개합니다.
팩토리 패턴(Factory Pattern)
팩토리 패턴은 객체 생성 책임을 한 곳에 위임하여 코드의 결합도를 낮추고 유지보수를 용이하게 하는 패턴입니다. 객체 생성 로직을 한 곳에서 관리하기 때문에 새로운 객체 타입이 추가되거나 기존 객체 생성 로직이 변경될 때 다른 코드에 미치는 영향을 최소화할 수 있습니다.
정합성 검사 잡에서는 제출하는 JSON 기반의 설정(Config) 값에 따라 Reader, Translator, Reporter 등 서로 다른 객체를 런타임 시점에 동적으로 생성해야 합니다. 이러한 상황에서 필요한 객체를 쉽게 생성하기 위해 다음과 같이 심플 팩토리 패턴을 활용했습니다.
아래는 데이터소스 타입에 따라 다른 Reader 객체를 반환하는 팩토리 코드입니다.
object DatasourceReaderFactory {
def getReader(
datasourceConfig: DatasourceConfig,
modeConfig: ModeConfig
): DatasourceReader = {
datasourceConfig match {
case datasourceConfig: MysqlDatasourceConfig => new MysqlReader(datasourceConfig, modeConfig)
case datasourceConfig: MongoDatasourceConfig => new MongoReader(datasourceConfig, modeConfig)
case datasourceConfig: IcebergDatasourceConfig => new IcebergReader(datasourceConfig, modeConfig)
case _ => throw new IllegalArgumentException(s"Unsupported datasource type: ${datasourceConfig.datasourceType}")
}
}
}
엔트리 클래스의 메인 함수에서는 위의 팩토리를 호출하여 JSON 설정(Config)에서 로드한 설정값을 전달하고 입력받은 설정 값에 따라 적절한 Reader 객체를 동적으로 생성합니다.
// 메인 함수 내 객체 생성 예시
val sourceReader = DatasourceReaderFactory.getReader(config.source, config.mode)
val targetReader = DatasourceReaderFactory.getReader(config.target, config.mode)
이러한 구조를 통해 새로운 데이터소스가 추가되거나 Reader 생성 조건이 변경되어도 팩토리 코드만 수정하면되어 유지보수와 확장성이 향상됩니다.
템플릿 메서드 패턴(Template Method Pattern)
템플릿 메서드 패턴은 여러 클래스에서 공통으로 사용하는 전체 로직의 뼈대(템플릿)를 상위 클래스에 정의하고, 세부적으로 달라지는 동작은 하위 클래스에서 구현하도록 하는 패턴입니다. 이를 통해 코드의 중복을 줄이고 공통 로직과 세부 구현 로직을 명확하게 분리할 수 있습니다.
정합성 검사 잡의 Reader 컴포넌트는 각 스캔 모드마다 서로 다른 방식으로 데이터를 읽습니다. 하지만 데이터를 읽은 이후 특정 컬럼 제외, 스키마 로깅 등 모든 모드에서 공통적으로 수행해야 하는 로직이 존재합니다. 이러한 로직을 쉽게 공통화하기 위해 템플릿 메서드 패턴을 활용했습니다.
아래는 각 Reader의 부모 클래스인 DatasourceReader의 read 함수에서 각 스캔 모드에 따라 데이터를 읽고 공통 로직을 처리하는 코드 예시입니다.
abstract class DatasourceReader(
datasourceConfig: DatasourceConfig,
modeConfig: ModeConfig
) {
def read(): DataFrameWrapper = {
// 각 모드에 따른 세부 동작 (하위 클래스에 구현)
val df = modeConfig.scanType match {
case "fullscan" => readWithFullScanMode()
case "keybased" => readWithKeyBasedMode()
case "limit" => readWithLimitMode()
case other => throw new IllegalArgumentException(s"Unsupported scan type: $other")
}
// 모든 read가 공통으로 처리해야하는 로직
val resultDF = df.drop(datasourceConfig.excludedColumns.toSeq: _*)
logger.info(s"DF schema: ${resultDF.schema}")
new DataFrameWrapper(resultDF, getOriginalColumnInfo)
}
// 하위 클래스에서 구현해야 하는 추상 메서드 정의
def readWithFullScanMode(): DataFrame
def readWithKeyBasedMode(): DataFrame
def readWithLimitMode(): DataFrame
이처럼 상위 클래스에서 스캔 로직의 큰 틀을 잡아두고, 세부 로직은 하위 클래스에 위임하는 방식으로 설계하면 코드의 중복을 효과적으로 줄일뿐 아니라 필수적인 공통 로직을 일관성 있게 수행할 수 있습니다.
전략 패턴(Strategy Pattern)
전략 패턴은 여러 알고리즘(행위)을 각각의 클래스로 캡슐화하고, 실행 시점에 적절한 알고리즘을 선택하여 유연하게 교체할 수 있도록 설계하는 패턴 입니다. 특히 동작 로직이 다양하고, 설정값에 따라 실행 동작이 달라져야 할 때 사용하면 좋은 패턴입니다. 앞서 소개한 Translator 컴포넌트는 전략 패턴의 전형적인 활용 사례입니다. JSON 포맷 보정, 해싱, 마스킹 등 다양한 변환 전략을 각각의 클래스로 정의하고, JSON 기반 config에 명시된 변환 타입 리스트에 따라 런타임 시 동적으로 Translator를 선택하여 적용합니다.
이처럼 Translator 인터페이스를 중심으로 전략을 분리하고 순차적으로 실행하는 구조를 통해 새로운 전략을 클래스 추가만으로 확장할 수 있고, 설정값을 통해 동작 방식을 유연하게 제어할 수 있습니다.
싱글톤 패턴(Singleton Pattern)
싱글톤 패턴은 객체의 인스턴스를 단 하나만 생성하도록 보장하는 패턴입니다. 프로그램 전체에서 반드시 하나의 인스턴스만 필요하거나, 하나의 인스턴스를 여러 곳에서 공유하고 싶을 때 주로 사용합니다. 특히 Scala에서는 object 키워드를 통해 클래스 수준에서 싱글톤 객체를 선언할 수 있어 별도의 구현 없이도 전역 인스턴스를 쉽게 관리할 수 있습니다.
정합성 검사 잡에서는 Spark를 사용해 데이터를 읽고 비교합니다. 이때 사용하는 Spark 세션은 무거워 전역적으로 공유되어야 하기 때문에, 싱글톤 객체로 관리하는 것이 가장 적합합니다. 특히 Production, Sandbox, Local 등 실행 환경에 따라 세션의 생성 방식이 다르기 때문에 생성 로직을 한 곳에 모아두고 관리하는 구조가 일관성과 유지보수 측면에서 매우 유용합니다. 아래는 싱글톤 객체로 구현한 SparkSessionBuilder의 예시 코드입니다.
import com.kakao.bi.consistency.config.AppConfig.stage
import org.apache.spark.sql.SparkSession
object SparkSessionBuilder {
val spark: SparkSession = stage match {
case "production" | "sandbox" =>
SparkSession.builder().getOrCreate()
case "local" =>
SparkSession
.builder()
.appName("syncverse-batch-consistency")
.master("local")
.getOrCreate()
case other => throw new IllegalArgumentException(s"Unsupported stage: $other")
}
}
이후 필요한 곳에서 아래와 같이 import 하여 동일한 인스턴스를 사용할 수 있습니다.
import com.kakao.bi.consistency.SparkSessionBuilder.spark
5. 코드 스타일 통일을 위한 자동화 방법
협업 과정에서는 코드 스타일의 일관성이 가독성과 유지보수성에 큰 영향을 미칩니다. 이 프로젝트에서는 다음 두 가지 도구를 활용하여 코드 스타일을 자동으로 관리하고 있습니다.
scalafmt
scalafmt은 Scala 코드를 지정된 스타일 규칙에 따라 자동으로 정렬해주는 도구입니다. 들여쓰기 크기, 최대 라인 길이, 주석 줄 바꿈 여부 등 여러 스타일에 대한 규칙을 .scalafmt.conf 설정 파일로 정의해 관리합니다.
version = "3.8.1"
runner.dialect = "scala212"
indent.main = 2
maxColumn = 120
docstrings.wrap = no
comments.wrap = no
scalafix
scalafix는 포맷팅뿐 아니라 불필요한 import 제거, import 순서 정렬 등 정적 분석 기반의 리팩토링 도구입니다. scalafmt만 사용할 경우 import 그룹 내에서의 순서는 유지되지만 그룹 간 순서는 정렬되지 않는데요. 이러한 문제를 해결하기 위해 scalafix를 함께 사용하여 보다 정교한 import 정렬을 수행합니다.
rules = [
OrganizeImports,
]
OrganizeImports {
groupedImports = Merge
}
이 두 도구는 CLI로 직접 실행할 수도 있지만 자동으로 적용될 수 있도록 Pre-commit hook에 연결해 커밋 시 자동으로 실행되도록 설정하여 사용 중입니다. 또한 GitHub에 코드를 푸시할 때도 내부 CI 훅을 통해 자동으로 포맷 검사를 수행하도록 구성해 일관된 코드 스타일이 유지되도록 하고 있습니다.
6. Part 1을 마무리하며
이번 코드를 설계하며 다양한 데이터 소스와 실행 환경에서도 유연하게 동작할 수 있는 정합성 검사 시스템을 만드는 데 집중했습니다. 역할별로 컴포넌트를 잘 나누고, 디자인 패턴을 적극적으로 도입하였고, 원활한 협업을 위한 스타일 포맷팅 통일까지 신경쓰는 등 결과적으로 확장성과 유지보수 측면에서 매우 만족스러운 형태로 구현할 수 있었습니다.
또한 무엇보다 단순히 동작만 잘하는 코드보다는 훗날 제가 없더라도 누구나 쉽게 이해하고, 수정하고, 개선해나갈 수 있는 튼튼하고 직관적인 구조를 만드는 것을 가장 중요한 철학으로 여겼습니다. 이전에는 여러 사람들이 작성한 코드가 각자 스타일도 다르고 정리도 안되어 있어서 파악과 유지보수가 어려운 경우가 많았는데요. 이번에는 보다 구조적인 설계를 통해 그런 어려움을 줄이고자 노력했습니다.
물론 아쉬운 부분도 남아 있습니다. 군데군데 Java 스타일의 문법이 남아 있어서 좀 더 Scala스럽게 간결하게 개선할 수 있었던 부분들이 있었고, Translator 구현 쪽도 돌아보면 "좀 더 나은 방법이 없었을까?"를 항상 고민하게 됩니다. 또한 enum을 쓰지 않고 문자열로 처리한 부분도 아쉬운 포인트로 남아있습니다.
그래도 이런 부족한 부분을 하나하나 인지하고 더 나은 방향으로 개선해 나가고자 늘 고민하고 있습니다. 전체적으로는 매우 재미있고 의미 있는 구조 설계 경험이었습니다.
다음 Part 2. Spark 최적화편에서는 Spark가 DB 데이터를 읽을 때 성능을 높이기 위해 적용한 여러 최적화 기법과 팁들을 자세히 소개할 예정입니다.