Python batch를 Spark로 옮겼더니 802.675가 달라졌다
Python으로 만든 batch transform을 Spark DataFrame으로 다시 쓰는 일은 문법 번역처럼 보인다.
Python loop / dict
-> Spark select / Window / groupBy
Enter fullscreen mode Exit fullscreen mode
하지만 실행 엔진만 바꿨는데도 결과가 달라질 수 있다. 중복 중 어떤 행을 남기는지, 평균을 어떻게 반올림하는지, quality 실패가 orchestration에 어떻게 전달되는지, 같은 입력을 정말 skip해도 되는지가 모두 계약이기 때문이다.
manufacturing-data-platform-mini의 S7에서는 Kafka landing에서 만든 한 날짜의 canonical CSV를 기존 Python batch와 새 Spark batch에 함께 넣었다. 목표는 “Spark를 사용했다”가 아니었다.
같은 입력을 다른 엔진으로 처리해도 grain, dedup, metric, quality, publish 상태 전이가 같아야 한다.
그 계약을 테스트하던 중 실제로 802.675 반올림 결과가 갈리는 false-green도 잡았다.
Scenario
앞선 slice에서는 local Kafka에서 받은 합성 설비 event를 immutable JSONL로 landing한 뒤, 한 business_date의 accepted event를 canonical CSV와 source_hash로 고정했다.
S7은 그 경계부터 시작한다.
Kafka raw landing
-> K1.5 canonical CSV + source_hash
-> Spark silver / gold
-> existing quality checks
-> Iceberg business_date partition
-> thin Airflow CLI wrapper
Enter fullscreen mode Exit fullscreen mode
운영 시나리오는 작다.
한 날짜를 Spark batch로 backfill한다.
같은 source 재실행은 새 snapshot을 만들지 않는다.
정정 source는 대상 날짜 partition만 교체한다.
quality가 실패하면 Iceberg current를 바꾸지 않는다.
Enter fullscreen mode Exit fullscreen mode
Spark가 Kafka JSONL을 다시 해석하게 만들지 않았다. 이미 provenance와 형식을 검증한 adapter의 CSV와 source_hash를 그대로 입력 계약으로 재사용했다.
Decision Pressure: 엔진보다 먼저 고정할 것
Python 구현과 Spark 구현이 “같다”고 하려면 무엇을 비교해야 할까?
이번 slice에서는 다섯 가지를 먼저 고정했다.
계약 고정한 내용 깨지면 생기는 문제 input identity canonical CSV +source_hash
서로 다른 입력을 같은 결과처럼 비교
grain / dedup
natural key와 Kafka coordinate 기준 first
합계 중복 또는 임의 행 선택
metric rounding
Python gold와 같은 소수 둘째 자리 결과
엔진별 metric 불일치
quality gate
기존 quality suite를 그대로 적용
Spark 경로만 다른 품질 기준 사용
publish state
same-source skip, correction overwrite
snapshot noise 또는 stale skip
이것이 코드보다 먼저 필요했다. groupBy가 실행된다는 사실만으로는 기존 batch와 같은 시스템이 되지 않는다.
Contract 1: 검증된 입력 경계를 재사용한다
S7은 K1.5 adapter가 만든 결과만 받는다.
input rows = canonical CSV
input version = SHA-256 source_hash
date scope = exactly one business_date
Enter fullscreen mode Exit fullscreen mode
Spark가 raw Kafka envelope을 직접 읽게 하면 adapter에서 검증한 manifest, accepted status, coordinate, provenance 계약을 우회한다. 엔진 교체 범위를 줄이기 위해 source boundary는 바꾸지 않았다.
또한 요청한 날짜와 CSV 내부 날짜가 다르면 SparkSession을 띄우거나 table을 건드리기 전에 실패시킨다.
Contract 2: dedup의 “첫 행”도 정의한다
silver natural key는 다음과 같다.
(work_order_id, machine_id, event_time)
Enter fullscreen mode Exit fullscreen mode
Python 구현은 canonical CSV에서 먼저 나온 행을 남긴다. 그런데 Spark DataFrame에는 별도의 순서를 주지 않으면 “첫 행”이 결정적이지 않다.
adapter가 CSV를 Kafka coordinate 순으로 쓰므로 Spark도 같은 순서를 명시했다.
dedup_order = Window.partitionBy(
"work_order_id", "machine_id", "event_time"
).orderBy(
F.col("kafka_topic"),
F.col("kafka_partition").cast("long"),
F.col("kafka_offset").cast("long"),
)
deduped = (
filtered.withColumn("_rn", F.row_number().over(dedup_order))
.filter(F.col("_rn") == 1)
.drop("_rn")
)
Enter fullscreen mode Exit fullscreen mode
같은 natural key를 가진 두 event를 넣는 테스트에서 Python과 Spark가 같은 행 하나를 남기는지 확인했다.
Contract 3: round 이름이 같아도 결과는 같지 않다
처음에는 Python round가 bankers’ rounding을 쓰므로 Spark bround와 같을 것이라고 생각했다. 일반 샘플은 통과했다.
반례는 cycle_time_ms 합계 32107, 행 수 40인 평균이었다.
수학적 평균: 32107 / 40 = 802.675
Python round(value, 2): 802.67
Spark bround(value, 2): 802.68
Enter fullscreen mode Exit fullscreen mode
이 차이는 “둘 다 half-even”이라는 이름만 비교해서는 잡히지 않았다. 실행 중인 double 값과 각 함수의 처리 경로까지 포함해야 했다.
40,400개의 이 프로젝트용 정수비 표본을 비교했을 때 bround는 Python 결과와 204건 달랐다. format_number로 둘째 자리까지 만든 뒤 grouping comma를 제거하고 double로 바꾸는 built-in 표현은 같은 표본에서 mismatch가 없었다.
def _round_like_python(col, scale: int):
return F.regexp_replace(
F.format_number(col, scale), ",", ""
).cast("double")
Enter fullscreen mode Exit fullscreen mode
그리고 32107 / 40을 별도 golden test로 남겼다.
중요한 경계도 있다.
이것은 이 gold metric과 40,400개 정수비 표본에서 확인한 bounded parity다. 모든 double과 모든 scale에서 Python
round와 보편적으로 같다는 주장은 아니다.
Contract 4: Spark용 quality를 새로 만들지 않는다
Spark용 quality suite를 별도로 구현하면 Python 경로와 조용히 달라질 수 있다. 그래서 S7은 Spark가 만든 silver/gold row를 driver로 collect한 뒤 기존 build_quality_checks를 그대로 적용했다.
Spark transform
-> materialized silver/gold rows
-> existing checks
- row count reconciliation
- unit / defect conservation
- numeric range
- schema and business-date checks
-> pass일 때만 Iceberg write
Enter fullscreen mode Exit fullscreen mode
이 선택은 local bounded slice에는 적합하지만, distributed Spark-native quality evaluation은 아니다.
quality가 실패하면 두 가지가 함께 성립해야 한다.
Iceberg snapshot을 만들지 않는다.
CLI가 non-zero로 종료되어 Airflow task도 실패한다.
Enter fullscreen mode Exit fullscreen mode
처음 구현은 첫 번째만 지키고 quality_failed JSON을 출력한 뒤 exit code 0으로 끝났다. 그러면 BashOperator는 데이터가 publish되지 않았는데도 task를 성공으로 볼 수 있다. review에서 이를 잡아 CLI가 SystemExit(1)로 종료되도록 고쳤다.
Contract 5: state 파일만 보고 skip하지 않는다
Iceberg publish는 overwritePartitions()를 쓴다.
gold_dataframe(spark, gold_rows) \
.writeTo("local.db.gold_daily_metrics") \
.overwritePartitions()
Enter fullscreen mode Exit fullscreen mode
정정 source가 가진 business_date partition만 교체하고, 다른 날짜 partition은 보존한다(같은 날짜 정정을 partition overwrite로 표현하는 이유는 B5에서 다뤘다). 같은 table + business_date + source_hash 재실행은 새 snapshot을 만들지 않는다.
여기에도 false skip이 있었다.
evidence state에는 이전 snapshot_id가 남아 있음
warehouse는 비워졌거나 새로 만들어짐
같은 source_hash가 다시 들어옴
Enter fullscreen mode Exit fullscreen mode
state 파일의 hash만 보면 “이미 처리했다”고 skip하지만 실제 table은 비어 있을 수 있다. 그래서 skip 조건에 snapshot history membership을 추가했다.
same_source = previous_state["source_hash"] == source_hash
snapshot_exists = previous_state["snapshot_id"] in existing_snapshot_ids
action = "skip" if same_source and snapshot_exists else "write"
Enter fullscreen mode Exit fullscreen mode
snapshot expiry나 GC로 과거 snapshot이 history에서 사라지면 같은 partition을 한 번 더 쓰는 correct-but-extra rewrite가 생길 수 있다. 이번 local slice에서는 snapshot expiry를 실행하지 않았고, 데이터가 없는 상태에서 잘못 skip하는 것보다 재작성을 선택했다.
Airflow에는 로직을 넣지 않았다
Airflow DAG는 위 CLI 하나를 호출하는 single-task wrapper다.
Airflow
-> validated CLI command
-> adapter
-> Spark transform
-> quality gate
-> Iceberg publish
Enter fullscreen mode Exit fullscreen mode
transform, quality, SparkSession, Iceberg write를 DAG body에 넣지 않았다. 로직은 일반 Python test와 Spark integration test로 검증하고, Airflow에서는 DAG import, command wiring, local dags test만 확인했다.
이는 production scheduler/executor나 운영 Airflow 검증이 아니다.
Evidence
독립 review 후 전체 검증 결과는 다음과 같다.
base Python environment: 90 passed, 14 skipped
Spark-visible environment: 99 passed, 5 skipped
S7 Spark integration: 14 passed
runtime state checks: 8/8 passed
isolated Airflow DagBag: 5 passed
local airflow dags test: DagRun success, task exit 0
Enter fullscreen mode Exit fullscreen mode
실제 local Iceberg 상태 전이도 확인했다.
다른 날짜 D2 publish
-> source A publish: snapshot 1 -> 2
-> 같은 source A retry: skipped, snapshot 그대로
-> correction source B: snapshot 2 -> 3
-> 대상 D1 units_produced=200으로 교체
-> D2 rows는 그대로 보존
Enter fullscreen mode Exit fullscreen mode
Spark gold의 groupBy 실행 계획에서 Exchange도 관찰했다. 이것은 shuffle이 발생한다는 local execution-plan 학습 evidence이지 성능이나 대규모 처리 성과는 아니다.
구현과 테스트:
- Spark machine-event batch code
- Engine parity and failure tests
- Engine-swap decision record
- Verification log
Limitations
이번 글에서 증명하지 않은 것은 명확하다.
production / cluster Spark
대규모 성능·throughput 개선
full bronze/silver/gold Spark-Iceberg pipeline
Spark Structured Streaming 또는 direct Kafka-to-Iceberg
distributed Spark-native quality evaluation
concurrent Iceberg writer correctness
Iceberg commit과 JSON evidence write의 하나의 transaction
production Airflow operation
모든 float 입력에 대한 Python/Spark 반올림 동치
Enter fullscreen mode Exit fullscreen mode
테이블도 local gold_daily_metrics 하나뿐이다. 이 결과를 production lakehouse 구축이나 대규모 Spark 운영 경험으로 확대하지 않는다.
정리
엔진 교체에서 먼저 비교해야 할 것은 코드 모양이 아니다.
같은 입력 identity인가?
같은 grain과 dedup 행을 선택하는가?
같은 metric을 만드는가?
같은 실패를 실패로 보고하는가?
같은 retry와 correction 상태 전이를 만드는가?
Enter fullscreen mode Exit fullscreen mode
이번 slice에서 가장 값진 결과는 Spark 코드 자체보다 false-green 네 개를 계약과 테스트로 바꾼 일이었다.
반올림 parity 반례
stale snapshot state의 false skip
quality 실패의 exit 0
adapter provenance의 미지속
Enter fullscreen mode Exit fullscreen mode
새 엔진을 붙였다는 말보다, 엔진이 바뀌어도 무엇이 같아야 하는지 설명하고 실패 반례로 증명하는 편이 더 강한 evidence가 된다.
답글 남기기