Kafka → NiFi → Kudu 로 직접 흘리던 파이프라인에서 Kudu 디스크 압박과 NiFi 부하, 원천 보존 요구가 겹쳐 중간에 HDFS 원본 적재 계층을 두기로 한 구조다. 파싱과 분기를 착지 단계에서 분리해, 원본을 먼저 안전하게 남기고 가공은 배치로 재현한다. Kafka retention 이 짧아 실물 보관이 필요하다는 판단이 출발점이었다.
Spark Structured Streaming 으로 Kafka value 를 파싱하지 않고 원본 문자열 그대로 HDFS 에 적재한다. offset 은 checkpoint 로 관리해 at-least-once 를 보장하고, 10분 트리거 마이크로배치에 zstd 압축을 건다.
/data/bronze/raw/<토픽명>/ingest_dt=YYYY-MM-DD/
파티션 키를 이벤트 시각이 아니라 착지일(ingest_dt)로 잡는다. 지연 도착 데이터가 과거 파티션을 다시 건드리지 않게 하기 위함이며, 이벤트 시각 기준 재배치는 Job B 가 담당한다.
CDP YARN 에는 토픽별로 독립 애플리케이션을 띄워 장애를 격리한다. 한 토픽의 문제가 다른 토픽 적재를 멈추지 않게 하는 것이 목적이고, 토픽별로 executor 수와 출력 파티션 수를 따로 산정한다.
HDFS 의 원본을 읽어 explode 한 뒤 기존 분기 로직을 재현하고, PK 기준으로 중복을 제거해 이벤트 날짜 파티션에 Parquet 으로 저장한다. spark.sql.sources.partitionOverwriteMode=dynamic 으로 대상 파티션만 덮어써 멱등성을 확보한다.
CDE Airflow 에서 1시간마다 YARN ResourceManager REST API 로 스트리밍 앱이 RUNNING 인지 확인하고, 아니면 메일로 알린다.
경량 NiFi 흐름(ConsumeKafkaRecord → MergeContent → PutHDFS)이나 Kafka Connect HDFS Sink 로도 같은 착지를 만들 수 있다. 어느 쪽을 쓰든 small file 이 쌓이지 않도록 128~256MB 를 목표로 병합하는 단계가 반드시 필요하다.
${KAFKA_TRUSTSTORE_PW} 형태로 환경변수 주입한다. 한 번이라도 평문 노출된 값은 로테이션한다.