it

Delta Lake Event Driven Design — ออกแบบ Data

delta lake event driven design
Delta Lake Event Driven Design — ออกแบบ Data

Delta Lake Event Driven

Delta Lake Event Driven Design — ออกแบบ Data

Delta Lake Event Driven ACID Transaction Change Data Feed Streaming Spark Merge Upsert Schema Evolution Time Travel S3 ADLS GCS

FeatureDelta LakeApache IcebergApache Hudi
ACID TransactionsYesYesYes
Change Data FeedCDF (Built-in)Incremental ReadCDC (Built-in)
StreamingSpark Structured StreamingFlinkSpark/Flink
Merge/UpsertMERGE INTOMERGE INTOUpsert
Time TravelVersion/TimestampSnapshotTimeline
OptimizeOPTIMIZE + Z-ORDERCompactionCompaction

Change Data Feed & Merge

Delta Lake Event Driven Design — ออกแบบ Data
# === Change Data Feed & Merge Operations ===



# Read Change Data Feed

# changes_df = spark.read \

#     .format("delta") \

#     .option("readChangeFeed", "true") \

#     .option("startingVersion", 5) \

#     .option("endingVersion", 10) \

#     .table("events")

#

# # CDF columns: _change_type, _commit_version, _commit_timestamp

# # _change_type: insert, update_preimage, update_postimage, delete

#

# inserts = changes_df.filter("_change_type = 'insert'")

# updates = changes_df.filter("_change_type = 'update_postimage'")

# deletes = changes_df.filter("_change_type = 'delete'")



# Merge (Upsert) - Exactly-once Processing

# deltaTable = DeltaTable.forPath(spark, "delta://silver/users")

# deltaTable.alias("target").merge(

#     source_df.alias("source"),

#     "target.user_id = source.user_id"

# ).whenMatchedUpdate(set={

#     "name": "source.name",

#     "email": "source.email",

#     "updated_at": "source.event_time"

# }).whenNotMatchedInsert(values={

#     "user_id": "source.user_id",

#     "name": "source.name",

#     "email": "source.email",

#     "created_at": "source.event_time",

#     "updated_at": "source.event_time"

# }).execute()



@dataclass

class MergePattern:

    pattern: str

    use_case: str

    match_condition: str

    matched_action: str

    not_matched_action: str



patterns = [

    MergePattern("SCD Type 1 (Overwrite)",

        "อัพเดทข้อมูลล่าสุด ไม่เก็บประวัติ",

        "target.id = source.id",

        "UPDATE SET * (Overwrite ทุก Column)",

        "INSERT * (Insert Row ใหม่)"),

    MergePattern("SCD Type 2 (History)",

        "เก็บประวัติทุก Version",

        "target.id = source.id AND target.is_current = true",

        "UPDATE SET is_current=false, end_date=now()",

        "INSERT with is_current=true, start_date=now()"),

    MergePattern("Deduplication",

        "ลบ Duplicate Event",

        "target.event_id = source.event_id",

        "ไม่ทำอะไร (Skip)",

        "INSERT * (Insert เฉพาะ Event ใหม่)"),

    MergePattern("Delete + Insert",

        "Replace Partition ด้วยข้อมูลใหม่",

        "target.date = source.date",

        "DELETE (ลบ Partition เก่า)",

        "INSERT * (Insert ข้อมูลใหม่ทั้ง Partition)"),

]



print("=== Merge Patterns ===")

for m in patterns:

    print(f"\n  [{m.pattern}] {m.use_case}")

    print(f"    Match: {m.match_condition}")

    print(f"    Matched: {m.matched_action}")

    print(f"    Not Matched: {m.not_matched_action}")

เคล็ดลับ

  • CDF: เปิด Change Data Feed ทุก Table สำหรับ Event-driven
  • Merge: ใช้ Merge แทน Overwrite สำหรับ Exactly-once
  • Z-Order: Z-Order ตาม Column ที่ Filter บ่อยที่สุด
  • Partition: Partition ตาม Date ไม่เกิน 10K Partitions
  • Vacuum: VACUUM ทุกสัปดาห์ Retain 7 วัน ลด Storage Cost

Delta Lake คืออะไร

Open Source Storage Layer ACID Transaction S3 ADLS GCS Spark Merge Time Travel CDF Schema Evolution Optimize Vacuum Parquet

อ่านเพิ่ม: Apache Kafka เจาะลึก สอน Kafka Streams, Connect, Schema Regi · อ่านเพิ่ม: PostgreSQL ขั้นสูง สอน Indexing, Query Optimization, Replica · อ่านเพิ่ม: Event-Driven Architecture คืออะไร? สอนออกแบบระบบ Event Sourc

XM Legend · เทรดเดอร์ & ผู้สอน Forex 13 ปี

ผู้ก่อตั้ง SiamCafe ตั้งแต่ปี 1997 · เทรดเดอร์สาย Forex มากกว่า 13 ปี ได้รับการยกย่องเป็น XM Legend · แบ่งปันความรู้ Forex, ไอที, AI และการเทรด จากประสบการณ์จริงในตลาดจริง