6분
적재 파이프라인 구현
ETL 파이프라인
적재 파이프라인 구현
ETL 파이프라인 > ETL 파이프라인
학습 목표
- Python + neo4j 드라이버로 배치 적재를 구현한다
- UNWIND 배치와 트랜잭션 경계를 설정한다
파이프라인 골격
Python 공식 드라이버로 구현합니다.
from neo4j import GraphDatabase
import pandas as pd, hashlib
URI, AUTH = "bolt://localhost:7687", ("neo4j", "portkg2026")
def call_id(source: str, decl_no: str) -> str:
h = hashlib.sha1(f"{source}:{decl_no}".encode()).hexdigest()[:12]
return f"urn:mrn-compat:portkg:call:{h}" # 결정적 ID (멱등성의 핵심)
LOAD_ARRIVAL = """
UNWIND $rows AS r
MERGE (v:Vessel {callSign: r.callSign}) // IMO 확보 시 imoNumber로 승격
SET v.name = r.name
MERGE (p:Port {unLocode: r.portCode})
MERGE (pc:PortCall {callId: r.callId})
MERGE (v)-[:MAKES_CALL]->(pc)
MERGE (pc)-[:AT_PORT]->(p)
MERGE (ae:ArrivalEvent {eventId: r.callId + ":arrival"})
SET ae.actualTime = datetime(r.arrivalUtc),
ae.timeBasis = "port", ae.source = "PORTMIS"
MERGE (ae)-[:PART_OF_CALL]->(pc)
"""
def load(df: pd.DataFrame, batch=1000):
rows = df.to_dict("records")
with GraphDatabase.driver(URI, auth=AUTH) as drv, drv.session() as s:
for i in range(0, len(rows), batch):
s.execute_write(lambda tx: tx.run(LOAD_ARRIVAL, rows=rows[i:i+batch]))
실전 요령 세 가지
① UNWIND 배치(1,000행 단위) — 행마다 쿼리 1개씩 날리면 100배 느립니다. 리스트를 통째로 넘기고 Cypher 안에서 풀어냅니다.
② 트랜잭션은 execute_write로 — 배치 단위 원자성을 확보합니다. ACID가 일하는 지점입니다.
③ 관계도 MERGE로 — CREATE면 재적재 때 관계가 중복됩니다. 노드만 멱등하게 만들고 관계를 CREATE로 두는 실수가 흔합니다.
실습
- 위 골격을 자기 매핑 시트에 맞게 고쳐 적재합니다.
- 같은 파일을 두 번 돌리고 노드·관계 수를 비교합니다. 숫자가 같으면 멱등성 확보입니다. 다르면 어느 MERGE가 키를 안 쓰고 있는지 찾으세요.
callSign대신name을 MERGE 키로 바꿔 한 번 돌려 보세요 — 중복 노드가 실제로 어떻게 생기는지 눈으로 봐 두면 오래 기억에 남습니다. (확인 후 원복)- 배치 크기를 1과 1000으로 바꿔 소요 시간을 재 봅니다.
에디터 로딩 중...