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

Delta Lake Event Driven

Delta Lake Event Driven ACID Transaction Change Data Feed Streaming Spark Merge Upsert Schema Evolution Time Travel S3 ADLS GCS
เนื้อหาเกี่ยวข้อง — แนะนำให้อ่าน Immutable OS Fedora CoreOS Cloud Migration
| Feature | Delta Lake | Apache Iceberg | Apache Hudi |
|---|---|---|---|
| ACID Transactions | Yes | Yes | Yes |
| Change Data Feed | CDF (Built-in) | Incremental Read | CDC (Built-in) |
| Streaming | Spark Structured Streaming | Flink | Spark/Flink |
| Merge/Upsert | MERGE INTO | MERGE INTO | Upsert |
| Time Travel | Version/Timestamp | Snapshot | Timeline |
| Optimize | OPTIMIZE + Z-ORDER | Compaction | Compaction |
Change Data Feed & Merge

# === 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
เนื้อหาเกี่ยวข้อง — แนะนำให้อ่าน non farm payroll 2022
อ่านเพิ่ม: Apache Kafka เจาะลึก สอน Kafka Streams, Connect, Schema Regi · อ่านเพิ่ม: PostgreSQL ขั้นสูง สอน Indexing, Query Optimization, Replica · อ่านเพิ่ม: Event-Driven Architecture คืออะไร? สอนออกแบบระบบ Event Sourc
แนะนำเพิ่มเติม — แหล่งความรู้ Forex iCafeForex
เนื้อหาเกี่ยวข้อง — แนะนำให้อ่าน functional programming concepts





