4.1 데이터 로딩이란?
로딩(Loading)은 ETL 파이프라인의 마지막 단계로, 변환된 데이터를 최종 목적지에 저장하는 과정입니다. 올바른 로딩 전략은 쿼리 성능, 비용, 유지보수성, 그리고 데이터 신선도에 큰 영향을 미칩니다.
로딩 전략을 선택할 때 고려해야 할 핵심 요소:
- 데이터 볼륨: 처리할 데이터의 크기
- 빈도: 얼마나 자주 로딩하는가
- 레이턴시: 실시간성 요구사항
- 비용: 스토리지 및 컴퓨팅 비용
- 복잡도: 구현 및 유지보수 복잡도
4.2 Full Load vs Incremental Load
데이터 로딩의 가장 기본적인 선택은 전체 로드와 증분 로드 사이의 결정입니다.
Full Load (전체 로드)
소스의 모든 데이터를 매번 다시 로드하는 방식:
-- 1단계: 기존 데이터 삭제
TRUNCATE TABLE target_table;
-- 2단계: 전체 데이터 로드
INSERT INTO target_table
SELECT * FROM source_table;
장점:
- 구현이 간단함
- 데이터 일관성 보장
- 삭제된 레코드 자동 반영
- 디버깅이 쉬움
단점:
- 대용량 데이터에서 비효율적
- 네트워크 대역폭 낭비
- 처리 시간 증가
- 비용 증가 (특히 클라우드)
적용 사례:
- 소규모 레퍼런스 테이블 (국가 코드, 제품 카테고리 등)
- 일별 스냅샷 테이블
- 데이터가 자주 변경되는 경우
Incremental Load (증분 로드)
변경된 데이터만 로드하는 방식:
-- 마지막 로드 이후 변경된 데이터만 가져오기
INSERT INTO target_table
SELECT *
FROM source_table
WHERE updated_at > (
SELECT MAX(updated_at) FROM target_table
);
장점:
- 처리 시간 대폭 감소
- 네트워크 대역폭 절약
- 비용 효율적
- 실시간에 가까운 데이터 동기화 가능
단점:
- 구현 복잡도 증가
- 타임스탬프 컬럼 필요
- 삭제 처리가 어려움
- 초기 Full Load 필요
적용 사례:
- 대용량 트랜잭션 테이블
- 로그 데이터
- 자주 업데이트되는 마스터 데이터
전략 비교표
| 항목 | Full Load | Incremental Load |
|---|---|---|
| 구현 복잡도 | 낮음 | 중간~높음 |
| 처리 시간 | 오래 걸림 | 빠름 |
| 비용 | 높음 | 낮음 |
| 데이터 일관성 | 보장됨 | 주의 필요 |
| 삭제 처리 | 자동 | 별도 구현 필요 |
| 필수 요구사항 | 없음 | 타임스탬프/CDC |
| 적합한 데이터 크기 | < 100만 행 | > 100만 행 |
4.3 Upsert 패턴
Upsert(Update + Insert)는 데이터가 존재하면 업데이트하고, 없으면 삽입하는 패턴입니다. 증분 로드에서 가장 많이 사용됩니다.
패턴 1: MERGE 문 (표준 SQL)
MERGE INTO dim_customers AS target
USING staging.customers AS source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN
UPDATE SET
name = source.name,
email = source.email,
phone = source.phone,
updated_at = CURRENT_TIMESTAMP()
WHEN NOT MATCHED THEN
INSERT (customer_id, name, email, phone, created_at, updated_at)
VALUES (
source.customer_id,
source.name,
source.email,
source.phone,
CURRENT_TIMESTAMP(),
CURRENT_TIMESTAMP()
);
패턴 2: INSERT ... ON CONFLICT (PostgreSQL)
INSERT INTO products (product_id, name, price, stock)
VALUES (101, 'MacBook Pro', 2499.99, 50)
ON CONFLICT (product_id)
DO UPDATE SET
name = EXCLUDED.name,
price = EXCLUDED.price,
stock = EXCLUDED.stock,
updated_at = CURRENT_TIMESTAMP();
패턴 3: REPLACE INTO (MySQL)
REPLACE INTO users (user_id, username, email)
VALUES (1, 'john_doe', 'john@example.com');
주의: REPLACE는 DELETE + INSERT이므로 AUTO_INCREMENT가 증가합니다.
패턴 4: Delta Lake MERGE (Spark)
from delta.tables import DeltaTable
deltaTable = DeltaTable.forPath(spark, "/data/events")
deltaTable.alias("target").merge(
source.alias("source"),
"target.event_id = source.event_id"
).whenMatchedUpdateAll(
).whenNotMatchedInsertAll(
).execute()
패턴 5: 두 단계 Upsert (호환성 우선)
-- 1단계: 업데이트
UPDATE target_table t
SET
name = s.name,
email = s.email,
updated_at = CURRENT_TIMESTAMP()
FROM staging_table s
WHERE t.user_id = s.user_id;
-- 2단계: 새 레코드 삽입
INSERT INTO target_table (user_id, name, email, created_at)
SELECT user_id, name, email, CURRENT_TIMESTAMP()
FROM staging_table s
WHERE NOT EXISTS (
SELECT 1 FROM target_table t WHERE t.user_id = s.user_id
);
Upsert 성능 최적화
- 인덱스: JOIN 키에 인덱스 생성
- 배치 크기: 10K-100K 레코드 단위로 처리
- 파티셔닝: 날짜별 파티션으로 스캔 범위 축소
- 병렬 처리: 파티션별 병렬 Upsert
4.4 벌크 로딩 (Bulk Loading)
대용량 데이터를 효율적으로 로드하는 기법입니다. 단건 INSERT보다 100배 이상 빠릅니다.
CSV 파일에서 벌크 로드
-- PostgreSQL COPY
COPY orders (order_id, user_id, total, order_date)
FROM '/tmp/orders.csv'
WITH (FORMAT csv, HEADER true);
-- MySQL LOAD DATA
LOAD DATA INFILE '/tmp/orders.csv'
INTO TABLE orders
FIELDS TERMINATED BY ','
ENCLOSED BY '"'
LINES TERMINATED BY '\n'
IGNORE 1 ROWS;
클라우드 스토리지에서 벌크 로드
Amazon Redshift
COPY orders
FROM 's3://my-bucket/data/orders/'
IAM_ROLE 'arn:aws:iam::123456789:role/RedshiftLoadRole'
FORMAT AS PARQUET;
Google BigQuery
bq load \
--source_format=PARQUET \
--autodetect \
my_dataset.orders \
gs://my-bucket/data/orders/*.parquet
Snowflake
COPY INTO orders
FROM @my_s3_stage/orders/
FILE_FORMAT = (TYPE = PARQUET)
MATCH_BY_COLUMN_NAME = CASE_INSENSITIVE;
Python에서 벌크 로드
import pandas as pd
from sqlalchemy import create_engine
# 데이터 준비
df = pd.read_csv('orders.csv')
# 데이터베이스 연결
engine = create_engine('postgresql://user:pass@localhost/db')
# 벌크 로드 (기본: 1000 rows씩 배치)
df.to_sql(
'orders',
engine,
if_exists='append',
index=False,
method='multi', # 멀티 로우 INSERT
chunksize=10000 # 10K씩 배치
)
Spark에서 벌크 로드
df = spark.read.parquet("s3://bucket/data/")
df.write \
.format("jdbc") \
.option("url", "jdbc:postgresql://localhost/db") \
.option("dbtable", "orders") \
.option("user", "user") \
.option("password", "password") \
.option("batchsize", 100000) \
.mode("append") \
.save()
벌크 로딩 최적화 팁
| 기법 | 설명 | 성능 향상 |
|---|---|---|
| 인덱스 비활성화 | 로드 전 인덱스 제거, 로드 후 재생성 | 2-5배 |
| 제약조건 비활성화 | FK, CHECK 제약 임시 비활성화 | 1.5-3배 |
| 압축 포맷 사용 | Parquet, ORC 사용 | 5-10배 |
| 병렬 로드 | 파티션별 병렬 로딩 | N배 (워커 수) |
| 배치 크기 조정 | 10K-100K rows per batch | 10-100배 |
4.5 스트리밍 로딩
실시간 데이터를 지속적으로 로드하는 방식입니다. 배치 로딩과 달리 데이터가 도착하는 즉시 처리합니다.
Kafka to Database
from kafka import KafkaConsumer
import psycopg2
consumer = KafkaConsumer(
'orders-topic',
bootstrap_servers=['localhost:9092'],
group_id='order-loader'
)
conn = psycopg2.connect("dbname=mydb user=user")
cursor = conn.cursor()
batch = []
BATCH_SIZE = 1000
for message in consumer:
order = json.loads(message.value)
batch.append((
order['order_id'],
order['user_id'],
order['total']
))
# 배치가 차면 벌크 삽입
if len(batch) >= BATCH_SIZE:
cursor.executemany(
"INSERT INTO orders VALUES (%s, %s, %s)",
batch
)
conn.commit()
batch = []
Spark Structured Streaming
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("StreamingLoader").getOrCreate()
# Kafka에서 스트림 읽기
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "orders") \
.load()
# JSON 파싱
orders = df.selectExpr("CAST(value AS STRING)") \
.select(from_json("value", schema).alias("data")) \
.select("data.*")
# 데이터베이스에 스트리밍 쓰기
query = orders.writeStream \
.foreachBatch(lambda df, epoch_id: df.write.jdbc(
url="jdbc:postgresql://localhost/db",
table="orders",
mode="append"
)) \
.trigger(processingTime='10 seconds') \
.start()
query.awaitTermination()
Flink to PostgreSQL
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream orders = env
.addSource(new FlinkKafkaConsumer<>("orders", schema, properties));
orders.addSink(JdbcSink.sink(
"INSERT INTO orders (order_id, user_id, total) VALUES (?, ?, ?)",
(statement, order) -> {
statement.setLong(1, order.getOrderId());
statement.setLong(2, order.getUserId());
statement.setDouble(3, order.getTotal());
},
JdbcExecutionOptions.builder()
.withBatchSize(1000)
.withBatchIntervalMs(200)
.withMaxRetries(3)
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:postgresql://localhost/db")
.withDriverName("org.postgresql.Driver")
.build()
));
env.execute("Stream to PostgreSQL");
스트리밍 로딩 고려사항
- 백프레셔: 소스가 목적지보다 빠를 때 처리
- 체크포인팅: 장애 복구를 위한 상태 저장
- 정확히 한 번: 중복 방지 메커니즘
- 지연 허용: 늦게 도착하는 데이터 처리
4.6 적재 목적지별 전략
Data Warehouse
| 플랫폼 | 권장 전략 | 파일 포맷 | 특징 |
|---|---|---|---|
| Snowflake | Bulk COPY | Parquet | 자동 최적화, 클러스터링 |
| BigQuery | Load Jobs | Avro, Parquet | 서버리스, 파티셔닝 필수 |
| Redshift | COPY from S3 | Parquet, ORC | DIST KEY, SORT KEY 중요 |
| Synapse | PolyBase | Parquet | 분산 처리, 통계 필요 |
Data Lake
# Parquet로 파티션 저장
df.write \
.partitionBy("year", "month", "day") \
.mode("append") \
.parquet("s3://lake/orders/")
# 파티션 구조:
# s3://lake/orders/year=2025/month=12/day=26/part-0000.parquet
OLTP Database
-- 트랜잭션 처리
BEGIN;
INSERT INTO orders VALUES (...);
INSERT INTO order_items VALUES (...);
UPDATE inventory SET stock = stock - 1 WHERE product_id = 123;
COMMIT;
4.7 파티셔닝 전략
날짜 파티셔닝 (가장 일반적)
-- PostgreSQL
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
total DECIMAL(10,2),
order_date DATE
) PARTITION BY RANGE (order_date);
CREATE TABLE orders_2025_01 PARTITION OF orders
FOR VALUES FROM ('2025-01-01') TO ('2025-02-01');
CREATE TABLE orders_2025_02 PARTITION OF orders
FOR VALUES FROM ('2025-02-01') TO ('2025-03-01');
해시 파티셔닝
CREATE TABLE users (
user_id BIGINT,
username VARCHAR(50),
email VARCHAR(100)
) PARTITION BY HASH (user_id);
CREATE TABLE users_p0 PARTITION OF users
FOR VALUES WITH (MODULUS 4, REMAINDER 0);
CREATE TABLE users_p1 PARTITION OF users
FOR VALUES WITH (MODULUS 4, REMAINDER 1);
-- ... p2, p3
리스트 파티셔닝
CREATE TABLE sales (
sale_id BIGINT,
region VARCHAR(20),
amount DECIMAL(10,2)
) PARTITION BY LIST (region);
CREATE TABLE sales_asia PARTITION OF sales
FOR VALUES IN ('KR', 'JP', 'CN', 'SG');
CREATE TABLE sales_europe PARTITION OF sales
FOR VALUES IN ('UK', 'FR', 'DE', 'IT');
4.8 인덱싱 전략
인덱스 타입별 활용
| 인덱스 타입 | 용도 | 예시 |
|---|---|---|
| B-Tree (기본) | 범위 검색, 정렬 | 날짜, ID, 숫자 |
| Hash | 등호 검색 | 고유 키 |
| GIN | 배열, JSON, 전문검색 | 태그, 문서 |
| GiST | 지리 데이터 | 위치 정보 |
| BRIN | 대용량 정렬 데이터 | 시계열 데이터 |
-- B-Tree 인덱스
CREATE INDEX idx_orders_date ON orders(order_date);
-- 복합 인덱스
CREATE INDEX idx_orders_user_date ON orders(user_id, order_date DESC);
-- 부분 인덱스
CREATE INDEX idx_active_users ON users(user_id)
WHERE status = 'active';
-- 표현식 인덱스
CREATE INDEX idx_lower_email ON users(LOWER(email));
4.9 성능 최적화
로딩 전 최적화
-- 인덱스 비활성화
DROP INDEX idx_orders_user_date;
DROP INDEX idx_orders_date;
-- 제약조건 비활성화
ALTER TABLE orders DISABLE TRIGGER ALL;
-- 벌크 로드 실행
COPY orders FROM '/data/orders.csv' CSV;
-- 인덱스 재생성
CREATE INDEX idx_orders_user_date ON orders(user_id, order_date);
CREATE INDEX idx_orders_date ON orders(order_date);
-- 제약조건 활성화
ALTER TABLE orders ENABLE TRIGGER ALL;
-- 통계 업데이트
ANALYZE orders;
병렬 로딩
import concurrent.futures
from sqlalchemy import create_engine
def load_partition(partition_date):
engine = create_engine(DATABASE_URL)
df = pd.read_parquet(f's3://bucket/data/{partition_date}/')
df.to_sql('orders', engine, if_exists='append', index=False)
print(f"Loaded {partition_date}")
dates = ['2025-12-01', '2025-12-02', '2025-12-03', ...]
with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:
executor.map(load_partition, dates)
압축 포맷 비교
| 포맷 | 압축률 | 읽기 속도 | 쓰기 속도 | 용도 |
|---|---|---|---|---|
| CSV | 1x | 느림 | 빠름 | 간단한 전송 |
| Parquet | 5-10x | 빠름 | 중간 | 분석 워크로드 |
| ORC | 5-10x | 빠름 | 중간 | Hive, Spark |
| Avro | 3-5x | 중간 | 빠름 | 스트리밍 |
| JSON | 1-2x | 느림 | 빠름 | API, 로그 |
4.10 에러 핸들링
배드 레코드 처리
-- Redshift: 에러 허용
COPY orders FROM 's3://bucket/data/'
IAM_ROLE 'arn:aws:iam::123456789:role/RedshiftRole'
MAXERROR 100; -- 100개까지 에러 허용
-- BigQuery: 에러 테이블 저장
bq load \
--max_bad_records=100 \
--error_destination_table=dataset.load_errors \
dataset.orders \
gs://bucket/data/*.csv
Dead Letter Queue 패턴
def load_with_dlq(records):
good_records = []
bad_records = []
for record in records:
try:
validated = validate_and_transform(record)
good_records.append(validated)
except Exception as e:
bad_records.append({
'record': record,
'error': str(e),
'timestamp': datetime.now()
})
# 정상 레코드는 타겟 테이블로
if good_records:
load_to_target(good_records)
# 에러 레코드는 DLQ로
if bad_records:
load_to_dlq(bad_records)
요약
이 장에서는 데이터 로딩의 다양한 전략과 기법을 배웠습니다:
- Full Load vs Incremental: 데이터 크기와 변경 빈도에 따른 선택
- Upsert 패턴: 데이터 병합을 위한 다양한 SQL 기법
- 벌크 로딩: 대용량 데이터를 효율적으로 로드하는 방법
- 스트리밍 로딩: 실시간 데이터 처리
- 파티셔닝과 인덱싱: 쿼리 성능 최적화
- 성능 최적화: 병렬 처리, 압축, 인덱스 관리
핵심 원칙
- 데이터 크기에 맞는 로딩 전략 선택
- 가능한 한 벌크 로딩 사용
- 파티셔닝으로 쿼리 성능 향상
- 압축 포맷으로 비용 절감
- 에러 핸들링으로 데이터 품질 보장
복습 문제
- Full Load와 Incremental Load의 차이점과 각각의 적용 사례를 설명하세요.
- Upsert 작업을 구현하는 3가지 이상의 방법을 제시하세요.
- 벌크 로딩이 단건 INSERT보다 빠른 이유는 무엇인가요?
- 1억 건의 데이터를 로드할 때 고려해야 할 최적화 기법을 나열하세요.
- 파티셔닝 전략 중 날짜 파티셔닝이 가장 많이 사용되는 이유는?
- 스트리밍 로딩과 배치 로딩의 차이점과 각각의 장단점을 비교하세요.
- Parquet 포맷이 CSV보다 분석 워크로드에 적합한 이유를 설명하세요.
- Dead Letter Queue 패턴이 필요한 이유와 구현 방법을 설명하세요.
실습 과제
- 100만 건의 CSV 데이터를 PostgreSQL에 로드하는 스크립트를 작성하세요 (벌크 로딩 사용).
- Incremental Load 파이프라인을 구현하세요 (타임스탬프 기반).
- Kafka에서 데이터를 읽어 데이터베이스에 스트리밍 로드하는 프로그램을 작성하세요.
- 날짜별 파티셔닝된 테이블을 생성하고, 쿼리 성능을 비교하세요.