블록체인 트랜잭션을 Spark로 처리하는 스트리밍 파이프라인
Kafka, Spark Structured Streaming, Delta Lake를 이용해 블록체인 트랜잭션을 수집·정규화하고 Reorg를 보정하는 설계 방법
2026-08-14 · 최초 발행 2025-11-09
블록체인 이벤트 분석에서 먼저 다뤄야 할 문제
블록체인 데이터는 블록, 트랜잭션, 이벤트 로그(Log), 상태(state)로 이루어진다. 불변 데이터라는 성격만으로 처리 방식을 결정할 수는 없다. 체인 재구성(Reorg), 최종성(finality) 지연, 블록 타임에 기반한 시계열 특성을 함께 고려해야 하기 때문이다.
Spark Structured Streaming은 Kafka, RPC, 파일 같은 소스에서 데이터를 받아 마이크로배치 또는 연속 처리로 연결한다. 상태 기반 연산, 워터마크, 사실상-정확히-한-번 처리语义를 사용할 수 있어 대용량 이벤트를 다루는 기반이 된다.
전체 흐름은 노드에서 데이터를 받아 정규화하고, 저장·서빙 계층으로 전달한 뒤 운영 거버넌스로 관리하는 구조다. 이때 이벤트 스키마(ABI), 재처리와 Reorg 관리, Delta·Apache Hudi·Iceberg 기반의 증분 저장, 카탈로그와 메타데이터 관리가 핵심이 된다.
수집부터 서빙까지 이어지는 계층
수집 계층은 풀 노드(Geth/Besu), 아카이브 노드, L2 시퀀서, Kafka 파이어호스에서 데이터를 받는다. 고가용성과 백프레셔 제어가 필요하며, ABI를 기준으로 로그를 디코딩하고 체인·네트워크·컨트랙트 버전을 태깅한다. 트랜잭션 키는 tx_hash+log_index로 표준화하고, 이벤트 시간은 block_timestamp로 보정한다.
Spark 처리 계층에서는 상태 기반 집계, 윈도우링, 조인, 이상치 탐지를 수행한다. 워터마크와 중복 제거는 정확성에 관여하고, 체크포인트와 트랜잭션 로그는 장애 복구에 쓰인다. 스케일 아웃 환경에서는 셔플 튜닝과 키 스큐 완화가 필요하다.
저장소는 Delta Lake, Hudi, Iceberg 같은 데이터 레이크 포맷을 활용해 증분 머지, 타임 트래블, 스키마 진화를 지원한다. 이후 데이터는 OLAP, 피쳐 스토어, 서빙 API로 나뉘며 BI, 리스크 엔진, 온체인 자동화 피드백 루프에 사용된다.
확정 상태를 다루기 위해서는 N-블록 컨펌을 기준으로 확정 데이터와 비확정 데이터 레이어를 분리한다. Reorg가 감지되면 MERGE나 DELETE로 데이터를 보정하고 오딧 로그를 남긴다. 이벤트 버전을 관리하면 하류 시스템에 미치는 영향을 줄일 수 있다.
키와 시크릿은 분리하고 노드 접근을 제어해야 한다. 데이터 민감도 분류, PII·메타데이터 마스킹, 접근 감사도 필요하다. ABI 카탈로그와 스키마 레지스트리를 운영하며, 변경 시에는 롤링 배포와 호환성 검증을 수행한다.
데이터가 이동하고 Reorg가 반영되는 흐름
입력 단계에서는 노드의 WebSocket에서 블록·트랜잭션·로그 스트림을 받고, 수집 에이전트가 체인ID와 블록번호를 기준으로 Kafka에 파티셔닝 생산한다. 파티션키는 tx_hash이며, 컨트랙트 시그니처는 ABI·스키마 레지스트리에서 조회한다.
Spark는 Kafka 이벤트를 마이크로배치로 수신해 이벤트 시간을 block_timestamp로 설정하고 워터마크를 적용한다. UDF로 로그를 디코딩한 뒤 tx_hash+log_index 기준으로 중복을 제거하고 상태 기반 윈도우 집계를 수행한다. Reorg가 발생하면 영향 블록 범위를 다시 처리해 Delta MERGE로 정정한다.
정규화 결과는 Delta Silver 레이어에, 집계와 서빙 결과는 Gold 레이어에 저장한다. 알림·리스크 API와 피쳐 스토어를 갱신하고, 필요한 경우 온체인 콜백을 트리거할 수 있다.
디코딩에 실패한 이벤트는 사선 옆 저장(Dead-letter) 후 샘플링 리플레이한다. 파티션 스큐에는 커스텀 파티셔너나 살팅 키를 적용하고, 다운스트림 장애는 배압 전파와 트랜잭션 로그 기반 재시도로 대응한다. Kafka idempotent producer와 Spark exactly-once sink(Delta)를 사용하며, 예를 들어 12블록 컨펌 임계치를 두고 확정 전 레코드에는 pending 플래그를 둘 수 있다. 최종성이 확보되면 상태를 고정하고 이전 버전은 아카이브한다.
Kafka에서 Delta까지의 정규화 예시
전제조건: Spark 3.4+, Python 3.10, Delta Lake 2.4+, Kafka 클러스터, web3.py 6+, ABI 레지스트리 API
실행 환경: spark-submit with --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.4.1, io.delta:delta-core_2.12:2.4.0
# main.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, expr, to_timestamp, window, udf
from pyspark.sql.types import StringType, StructType, StructField, LongType
# 단순화한 ABI 디코더 (실무는 캐시+외부 레지스트리 사용)
def decode_log(topic0, data):
# topic0: 이벤트 시그니처, data: hex payload
# TODO: 실제 ABI 파싱 적용
return "Transfer" if topic0 == "0xddf252ad..." else "Unknown"
decode_udf = udf(decode_log, StringType())
spark = (SparkSession.builder
.appName("bc-spark-pipeline")
.config("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
.getOrCreate())
raw = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "eth-logs")
.option("startingOffsets", "latest")
.load())
schema = StructType([
StructField("chain_id", StringType()),
StructField("block_number", LongType()),
StructField("block_timestamp", LongType()),
StructField("tx_hash", StringType()),
StructField("log_index", LongType()),
StructField("topic0", StringType()),
StructField("data", StringType())
])
json = raw.selectExpr("CAST(value AS STRING) as v") \
.select(from_json(col("v"), schema).alias("j")) \
.select("j.*")
parsed = (json
.withColumn("event_name", decode_udf(col("topic0"), col("data")))
.withColumn("event_time", to_timestamp((col("block_timestamp")).cast("timestamp")))
.withWatermark("event_time", "10 minutes")
.dropDuplicates(["tx_hash", "log_index"]))
tps = (parsed
.groupBy(window(col("event_time"), "1 minute"), col("event_name"))
.count()
.withColumnRenamed("count", "events_per_minute"))
# Delta Sink
query = (tps.writeStream
.format("delta")
.option("checkpointLocation", "s3://bucket/chk/tps")
.outputMode("complete")
.start("s3://bucket/delta/gold/tps"))
query.awaitTermination()
이 파이프라인은 워터마크와 dropDuplicates로 재전송·중복을 막고, RocksDB state store로 대규모 상태를 안정화한다. Delta에는 complete 모드로 저장하며 체크포인트로 정확히-한-번 처리를 보장한다.
Reorg 영향 범위를 다시 처리한 뒤 Silver 테이블을 보정하는 코드는 다음과 같다.
// 영향 블록 범위 재처리 후 Silver 테이블 보정
spark.sql("""
MERGE INTO delta.`s3://bucket/delta/silver/logs` AS t
USING delta.`s3://bucket/tmp/reorg_patch` AS s
ON t.tx_hash = s.tx_hash AND t.log_index = s.log_index
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
이 방식은 재구성 범위의 레코드를 업서트하고, 타임 트래블을 이용한 전후 비교와 감사 로그 유지를 뒷받침한다.
수집 경로별 선택 기준
| 옵션 | 성능 | 확장성 | 일관성 | 안정성 | 운영 편의 |
|---|---|---|---|---|---|
| RPC 직접 폴링 → Spark | 중 | 중 | 중 | 중 | 보통 |
| Kafka 파이어호스(에이전트) → Spark | 상 | 상 | 상 | 상 | 좋음 |
| 관리형 인덱싱 서비스(The Graph 등) → 배치 | 중 | 상 | 중 | 상 | 매우 좋음 |
| 노드 DB 덤프 → 배치 ETL | 중 | 중 | 상 | 상 | 보통 |
실시간 대용량 처리와 재시도 제어가 필요하다면 Kafka를 경유하는 구성이 맞다. 강한 최종성과 감사가 요구되는 환경에서는 Delta나 Hudi를 활용한 업서트 패턴을 채택한다.
분석 결과가 쓰이는 곳
거래소 리스크·AML 모니터링에서는 온체인 이벤트에 지갑 군집화와 블랙리스트 조인을 적용하고 이상 전송 패턴을 탐지한다. 1~5분 SLA의 실시간 경보를 발행하고 자동 인출 제한으로 연결할 수 있다.
DeFi·NFT 운영 분석에서는 풀 유동성 변동, 스왑 슬리피지, 민팅 활동을 집계한다. 파생 지표는 피쳐 스토어에 공급해 마케팅과 수수료 정책 최적화에 사용한다.
공급망·원산지 추적에서는 온체인 자산 이력과 오프체인 IoT 텔레메트리를 결합하고, 이벤트 시간 정렬을 통해 프로세스 병목을 찾는다. CBDC·결제 인프라에서는 TPS와 확정 지연을 모니터링하며 결제 실패 원인을 분석하고 규제 보고용 데이터 품질을 보증한다.
처리 성능과 운영 조건
Kafka+Spark 구성은 클러스터와 키 스큐에 따라 초당 수만수십만 이벤트를 처리할 수 있으며, 정확성 우선 설정에서는 엔드투엔드 12분 SLA 달성이 가능하다. 최신 정보 확인이 필요하다.
레이크하우스를 채택하면 컬럼형 저장, ZOrder, 압축을 통해 스토리지 비용을 절감할 수 있다. 증분 처리와 자동 스케일링은 운영 인건비 절감에 연결된다. 워터마크, 중복 제거, 컨펌 레벨은 정확성을 높이고, 스키마 진화와 타임 트래블은 변경 대응과 감사를 수월하게 한다.
파티셔닝 키는 chain_id+block_number 모듈로 두어 병렬성을 높일 수 있다. 살팅, AQE, skew join은 키 스큐를 완화하지만 조인 비용을 높이는 트레이드오프가 있다. 컨펌 레벨을 높이면 오탐과 정정은 줄어드는 대신 지연이 커지며, Delta MERGE 주기는 비용과 일관성 사이에서 조정해야 한다.
노드 API 인증, 레이트 리밋, 네트워크 격리와 함께 HSM/KMS 기반 키 관리를 적용하고 비서명 데이터만 파이프라인으로 처리하는 방식을 권장한다. 성능보다 보안을 강화하면 지연은 증가할 수 있다. 체크포인트와 Prometheus, Ganglia 모니터링을 갖추고, 스키마나 ABI 변경은 카나리 배포로 검증한다.
운영 전 확인할 항목
- 수집 경로에 WebSocket+Kafka를 두고 재시도와 배압 전략을 수립한다.
- ABI 레지스트리와 Backward/Full 호환성 규칙으로 스키마 변경을 관리한다.
- 워터마크, 컨펌 레벨, Reorg 보정 프로시저를 일관성 정책에 포함한다.
- Delta 포맷, 파티션 설계, ZOrdering을 저장소 전략으로 검토한다.
- 지연, 처리량, 누락율, 재처리 카운터를 대시보드에서 관찰한다.