빅데이터 처리 파이프라인: 수집부터 스트림 분석·저장·시각화까지

Flume, Sqoop, Kafka, Spark, Storm, HDFS, Hive를 연결해 빅데이터 수집·처리·저장·분석 파이프라인을 운영하는 구조와 구현 예시

2026-08-14 · 최초 발행 2025-10-14

데이터 흐름이 멈추지 않게 만드는 처리 체계

빅데이터 처리는 대량·다종·고속 데이터를 안정적으로 받아 내고, 장애에 견디는 저장소에 보관한 뒤 배치 또는 스트림으로 처리해 분석과 시각화까지 연결하는 종단 간 파이프라인이다. 데이터가 늘고 실시간 분석 요구가 커질수록 수집→정제→적재→분석/시각화→감시의 흐름을 분리하면서도 끊기지 않게 운영해야 한다.

이 구성에서는 Flume, Sqoop, Kafka, Spark, Storm, HDFS, HBase, MapReduce, Hive, Pig, R, Python, Tableau, ZooKeeper가 각 계층의 역할을 나눠 맡는다.

수집 단계에서는 Flume이 로그와 파일 스트림을 받고, 채널 기반 백프레셔와 재전송을 통해 HDFS 또는 Kafka와 연결한다. 관계형 DB 데이터를 대량으로 옮길 때는 Sqoop을 사용하며, 병렬 맵 태스크와 증분 로드로 HDFS나 Hive 적재를 수행한다.

Kafka는 파티션과 오프셋을 바탕으로 확장되는 메시지 버스다. 적어도 한 번 또는 정확히 한 번 처리 구성을 선택할 수 있으며, Avro/JSON 스키마와 스키마 레지스트리를 함께 두는 구성이 권장된다. ZooKeeper는 Kafka, HBase, Storm의 메타데이터 및 리더 선출을 관리하고 세션·워처 기반 장애 감지를 지원한다.

처리 엔진의 선택은 지연 시간과 작업 방식에 달려 있다. Spark는 메모리 중심 분산 처리와 Structured Streaming을 통해 이벤트 타임, 윈도우, 체크포인트를 다루며 배치와 스트림 코드의 통합성을 제공한다. Storm은 튜플 단위의 저지연 처리와 토폴로지 기반 확장에 맞고, at-least-once를 기본으로 보장한다. MapReduce는 대용량 배치에 안정적이지만 레이턴시와 개발 생산성 측면에서 불리하므로 레거시 배치 유지에 적합하다.

저장 및 질의 계층도 접근 패턴에 따라 분리한다. HDFS는 블록 복제를 사용하는 스루풋 지향 분산 파일 저장소로, 스냅샷과 HA 네임노드 구성을 활용한다. HBase는 랜덤 읽기·쓰기와 와이드 컬럼 설계에 적합하지만 RegionServer 확장, WAL, Compaction 운용이 필요하다. Hive와 Pig는 스키마 온 리드 기반으로 동작하며, 파티셔닝·버킷팅·ORC/Parquet 컬럼 포맷을 활용할 수 있다. Hive ACID 테이블은 업데이트와 머지 요구를 다룬다.

분석 결과는 R, Python, Tableau로 이어진다. 이때 오프셋 지연, 처리율, 에러율, GC와 자원 사용량을 수집하고 알람 기준선을 세워야 한다. Kerberos 인증, Kafka ACL/HDFS Ranger 정책, 데이터 마스킹과 암호화도 운영 계층에 포함된다. Kafka idempotent producer, Spark 체크포인트와 트랜잭셔널 싱크, Hive ACID 테이블은 일관성과 트랜잭션 요구를 다루는 수단이다.

계층을 연결하는 데이터 아키텍처

분석/시각화저장/쿼리처리메시징/스트림수집파일/이벤트SinkRDBHDFS/Hive LoadETL/BatchServingStream UpdateLegacy Batch실패시 재시도/체크포인트실패 튜플 재처리실패 이벤트→DLQ앱/로그/DBFlumeSqoopKafkaZooKeeperSpark (Batch/Streaming)StormMapReduceHDFSHBaseHiveRPythonTableau

복구 경로도 데이터 경로와 함께 설계한다. Flume에서 실패한 이벤트는 DLQ로 보내고, Kafka 프로듀서는 재시도와 idempotence 설정을 사용한다. Spark는 체크포인트를 기준으로 재처리하며, Storm은 실패한 튜플을 재전송하는 메커니즘을 적용한다.

파이프라인이 쓰이는 데이터 흐름

클릭스트림 분석에서는 Flume→Kafka→Spark Structured Streaming→HDFS/Hive→Tableau 경로를 사용할 수 있다. 이벤트 타임 윈도우와 워터마크로 지연 데이터를 허용하고, ORC 파티션 테이블로 비용을 줄인다.

관계형 DB를 데이터 레이크로 옮겨 배치 집계할 때는 Sqoop→HDFS→Hive/MapReduce 또는 Spark 흐름을 둔다. 증분 키 기반의 일일 증분 로드, Hive ACID Merge를 이용한 업서트, 메타데이터 카탈로그 일원화가 이 경로의 운영 요소다.

이상 징후 탐지와 실시간 알림은 Kafka→Storm 또는 Spark→HBase 서빙→Python API→알림으로 구성할 수 있다. 상태 기반 윈도우 집계, ZooKeeper를 통한 토폴로지 재밸런싱, HBase 핫 리전 회피 설계가 함께 필요하다.

IoT 시계열 데이터는 게이트웨이→Flume/Kafka→Spark→HDFS 파케이→Tableau 경로로 집계한다. 장치ID와 시간을 기준으로 파티셔닝하고, Late arrival 처리와 스냅샷·롤업 테이블 운영을 고려한다.

처리 성능과 운영상 기대할 수 있는 변화

Kafka 파티션을 증설하면 선형 확장이 가능하며, 파티션당 수 MB/s 수준 처리가 가능하다(환경 의존). Spark와 Storm을 이용한 실시간 처리는 수초~수백ms 단위 SLA를 목표로 할 수 있다.

HDFS와 객체 스토리지를 계층화하면 저장 비용을 30% 이상 절감할 수 있다. 체크포인트, 재처리, 복제를 조합하면 데이터 손실 확률을 낮추고 운영 MTTD/MTTR 단축에도 연결된다.

Spark·Storm·MapReduce를 선택하는 기준

엔진 성능(지연/처리량) 확장성 일관성/처리 보장 안정성 운영 편의
Spark 낮은 지연~배치, 높은 처리량 수평 확장 용이 정확히 한 번 구성 가능(구성 필요) 성숙 코드 통합성 높음
Storm 매우 낮은 지연, 중간 처리량 수평 확장 적어도 한 번 기본, 정확히 한 번은 복잡 안정 토폴로지 운용 필요
MapReduce 높은 지연, 대용량 배치 적합 수평 확장 배치 정확성 우수 매우 안정 개발 생산성 낮음

수집·적재·스트리밍 구현 예시

전제조건: Hadoop 3.x, Spark 3.x, Kafka 3.x, ZooKeeper 3.8, Kerberos 선택 적용, Schema Registry 선택

Flume에서 Kafka로 로그 수집

# flume-ng agent -n a1 -f flume.conf
a1.sources = r1
a1.channels = c1
a1.sinks = k1

a1.sources.r1.type = exec
a1.sources.r1.command = tail -F /var/log/app/app.log
a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /var/lib/flume/checkpoint
a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink
a1.sinks.k1.topic = app_logs
a1.sinks.k1.brokerList = kafka-1:9092,kafka-2:9092
a1.sinks.k1.requiredAcks = 1

a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

Sqoop으로 Hive에 적재

sqoop import \
  --connect jdbc:mysql://db:3306/sales \
  --username app --password **** \
  --table orders \
  --hive-import --create-hive-table \
  --hive-table dw.orders_raw \
  --fields-terminated-by '\t' \
  --num-mappers 4 \
  --incremental append --check-column id --last-value 0

Kafka 데이터를 Parquet/Hive로 쓰는 Spark Structured Streaming

# spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.1 stream.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window
from pyspark.sql.types import StructType, StringType, TimestampType

spark = SparkSession.builder.appName("kafka_to_parquet").enableHiveSupport().getOrCreate()
schema = StructType().add("userId", StringType()).add("event", StringType()).add("ts", TimestampType())

df = (spark.readStream.format("kafka")
      .option("kafka.bootstrap.servers", "kafka-1:9092,kafka-2:9092")
      .option("subscribe", "app_logs")
      .option("startingOffsets", "latest")
      .load())

parsed = (df.select(from_json(col("value").cast("string"), schema).alias("j"))
            .select("j.*")
            .withWatermark("ts", "10 minutes"))

agg = (parsed.groupBy(window(col("ts"), "5 minutes"), col("event"))
              .count())

query = (agg.writeStream
            .outputMode("append")
            .format("parquet")
            .option("path", "hdfs:///data/app_logs/agg")
            .option("checkpointLocation", "hdfs:///chk/app_logs/agg")
            .start())
query.awaitTermination()

Hive 외부 테이블

CREATE EXTERNAL TABLE IF NOT EXISTS dw.app_logs_agg (
  window_start timestamp,
  window_end   timestamp,
  event        string,
  cnt          bigint
)
PARTITIONED BY (dt string)
STORED AS PARQUET
LOCATION 'hdfs:///data/app_logs/agg';

스키마·복구·보안에서 결정되는 운영 품질

Avro/Protobuf와 Schema Registry를 적용하고 하위 호환 스키마 진화 규칙을 지킨다. 정확히 한 번 처리가 필요한 경우 Kafka idempotent producer와 transactional.id를 사용하고, Spark 트랜잭셔널 싱크(Hudi/Iceberg/Delta)를 고려한다.

파티션 수와 컨슈머 동시성의 정합성을 유지하며, HDFS 블록 크기와 Hive 파티션을 균형 있게 구성한다. HBase는 프리스플릿으로 핫스팟을 막는다. DLQ 트픽을 운영하고 재처리 가능한 원본을 원시 존에 보관하며, 체크포인트와 스냅샷 주기를 관리한다.

보안과 감사에는 Kerberos/SPNEGO, Kafka/HDFS Ranger 정책·감사 로그, 민감 데이터 마스킹·암호화를 적용한다. Sqoop는 Apache Attic 이관(최신 정보 확인 필요), Kafka는 KRaft 모드로 ZooKeeper 비의존 운영 가능(최신 정보 확인 필요)하므로 대체 기술 검토 필요성을 인지해야 한다.

단계적 도입은 파일럿에서 시작해 점진적으로 확대한다. 스키마와 데이터 거버넌스를 먼저 갖추고, 정확히 한 번 처리와 보안 정책을 바탕으로 SLA 중심의 운영 체계를 구성한다.

빅데이터데이터 파이프라인스트림 처리카프카스파크하둡