# CDF로 SCD Type 2 구현 (foreachBatch 패턴)
def upsert_to_scd2(batch_df, batch_id):
from delta.tables import DeltaTable
# UPDATE: 기존 행의 end_date를 설정 (이력 닫기)
updates = batch_df.filter("_change_type = 'update_postimage'")
target = DeltaTable.forName(spark, "catalog.schema.customers_scd2")
target.alias("t").merge(
updates.alias("s"),
"t.customer_id = s.customer_id AND t.is_current = true"
).whenMatchedUpdate(set={
"is_current": "false",
"end_date": "s._commit_timestamp"
}).whenNotMatchedInsert(values={
"customer_id": "s.customer_id",
"name": "s.name",
"email": "s.email",
"start_date": "s._commit_timestamp",
"end_date": "null",
"is_current": "true"
}).execute()
# 스트림 실행
(changes_df
.writeStream
.foreachBatch(upsert_to_scd2)
.option("checkpointLocation", "/checkpoints/scd2")
.trigger(availableNow=True)
.start()
)