Data Engineering

تسلط بر دلتا لیک: چارچوب ذخیره‌سازی متن‌باز برای دریاچه‌های داده قابل اعتماد

مقدمه

دریاچه‌های داده سنتی مدت‌هاست که از مشکل «خندق داده» رنج می‌برند. اگرچه آن‌ها مقیاس‌پذیری و مقرون‌به‌صرفه بودن ذخیره‌سازی شیء (مانند AWS S3 یا Azure Blob Storage) را ارائه می‌دهند، اما اغلب فاقد قابلیت اطمینان، سازگاری و عملکردی هستند که از یک پایگاه داده انتظار می‌رود. دلتا لیک (Delta Lake) به عنوان یک لایه ذخیره‌سازی متن‌باز که قابلیت اطمینان و عملکرد را به دریاچه‌های داده می‌آورد، وارد میدان می‌شود. با ترکیب بهترین ویژگی‌های یک دریاچه داده با قابلیت‌های مدیریت یک انبار داده، دلتا لیک معماری دریاچه-انبار داده (Data Lakehouse) را ممکن می‌سازد و به استاندارد پایپ‌لاین‌های مهندسی داده مدرن تبدیل شده است.

قابلیت‌های اصلی دلتا لیک

دلتا لیک قابلیت‌های فایل‌های Parquet را با ارائه مدیریت متادیتا روی قالب‌های استاندارد فایل داده گسترش می‌دهد. این سیستم بر پایه یک لاگ تراکنش ساخته شده است که یکپارچگی داده را تضمین می‌کند. مهم‌ترین ویژگی‌ها برای هر مهندس داده عبارتند از:

  • تراکنش‌های ACID: تضمین می‌کند که داده‌ها حتی در طول خواندن و نوشتن همزمان نیز سازگار باقی بمانند. این ویژگی از خواندن‌های ناقص جلوگیری کرده و اطمینان حاصل می‌کند که در صورت شکست عملیات نوشتن، مجموعه داده در حالت اولیه خود باقی می‌ماند.
  • اجرای طرحواره و تکامل آن: دلتا لیک با اعمال طرحواره‌ها، از ورود داده‌های نامناسب به دریاچه داده جلوگیری می‌کند. همچنین امکان تکامل طرحواره را فراهم می‌سازد و به شما اجازه می‌دهد ستون‌های جدیدی اضافه کنید یا انواع داده را تغییر دهید بدون اینکه پایپ‌لاین‌های موجود را مختل کنید.
  • سفر در زمان (نسخه‌بندی): این ویژگی شاید قدرتمندترین قابلیت باشد. شما می‌توانید با مشخص کردن یک شماره نسخه یا یک مهر زمانی، به نسخه‌های قبلی داده‌های خود دسترسی پیدا کنید که امکان اشکال‌زدایی، حسابرسی و قابلیت بازگشت به عقب (Rollback) را به سادگی فراهم می‌کند.
  • یکپارچه‌سازی پردازش دسته‌ای و جریان‌یافته: دلتا لیک یک API یکسان برای هر دو پردازش دسته‌ای و جریان‌یافته (Streaming) ارائه می‌دهد که پیچیدگی نگهداری سیستم‌های جداگانه برای داده‌های تاریخی و بلادرنگ را ساده می‌سازد.

پیاده‌سازی دلتا لیک با PySpark

برای بهره‌برداری از دلتا لیک، معمولاً آن را همراه با Apache Spark استفاده می‌کنید. در زیر یک مثال عملی آورده شده است که نحوه ایجاد یک جدول دلتا، انجام عملیات نوشتن با اجرای طرحواره و استفاده از سفر در زمان برای پرس‌وجوی نسخه‌های قبلی را نشان می‌دهد.

from delta.tables import DeltaTable
from pyspark.sql import SparkSession

# راه‌اندازی جلسه Spark
spark = SparkSession.builder \
    .appName("DeltaLakeExample") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 1. ایجاد یک جدول دلتا با اجرای طرحواره
data = [
    ("1", "Alice", 30),
    ("2", "Bob", 25)
]
df = spark.createDataFrame(data, ["id", "name", "age"])

# ذخیره به عنوان یک جدول دلتا. این کار به طور خودکار لاگ تراکنش را ایجاد می‌کند.
delta_path = "/path/to/delta/table"
df.write.format("delta").mode("overwrite").save(delta_path)

# 2. سفر در زمان: پرس‌وجوی نسخه 0
# این کار داده‌ها را همان‌طور که در نسخه قبلی وجود داشتند بازیابی می‌کند
df_v0 = spark.read.format("delta").option("versionAsOf", 0).load(delta_path)
df_v0.show()

# 3. عملیات آپسرت (Upsert) یا ادغام (Merge)
new_data = [
    ("1", "Alice", 31),
    ("3", "Charlie", 35)
]
df_new = spark.createDataFrame(new_data, ["id", "name", "age"])

# ادغام داده‌های جدید با جدول موجود
delta_table = DeltaTable.forPath(spark, delta_path)
delta_table.alias("old").merge(
    df_new.alias("new"),
    "old.id = new.id"
).whenMatchedUpdateAll() \
 .whenNotMatchedInsertAll() \
 .execute()

در کد بالا، به استفاده از option("versionAsOf", 0) توجه کنید. این گزینه به شما امکان می‌دهد وضعیت دقیق داده‌ها را قبل از هرگونه تغییر پرس‌وجو کنید. علاوه بر این، تابع merge نشان می‌دهد که دلتا لیک عملیات آپسرت (Upsert) را که در دریاچه‌های داده مبتنی بر Parquet استاندارد، به طور مشهور دشوار و کند هستند، ساده می‌سازد.

بهینه‌سازی عملکرد: Vacuum و Optimize

از آنجا که دلتا لیک تاریخچه تغییرات را برای سفر در زمان و انطباق با ACID نگه می‌دارد، فایل‌های کوچک ممکن است در طول زمان انباشته شوند. برای حفظ عملکرد، دو دستور ضروری هستند:

  1. .optimize(): این دستور فایل‌های کوچک را فشرده کرده و به فایل‌های بزرگ‌تر و بهینه‌شده تبدیل می‌کند که سرعت پرس‌وجوهای خواندن را به طور قابل توجهی افزایش می‌دهد.
  2. .vacuum(): این دستور فایل‌های قدیمی که دیگر توسط لاگ تراکنش ارجاع داده نمی‌شوند را حذف می‌کند. هشدار: اطمینان حاصل کنید که سیاست نگهداری (Retention Policy) شما به درستی تنظیم شده است، زیرا پاکسازی (Vacuum) داده‌های تاریخی مورد نیاز برای سفر در زمان را که فراتر از دوره نگهداری هستند، به طور دائمی حذف می‌کند.

نتیجه‌گیری

دلتا لیک به طور مؤثر شکاف بین انعطاف‌پذیری دریاچه‌های داده و قابلیت اطمینان انبارهای داده را پر کرده است. برای مهندسان داده، این ابزار یک راه‌حل قدرتمند برای مشکلات رایجی مانند خرابی داده‌ها، انحراف طرحواره (Schema Drift) و عملکرد کند پرس‌وجو ارائه می‌دهد. با یکپارچه‌سازی تراکنش‌های ACID، سفر در زمان و پردازش یکپارچه در یک چارچوب متن‌باز، دلتا لیک زیربنایی برای ساخت معماری‌های داده مقیاس‌پذیر، قابل اعتماد و با عملکرد بالا فراهم می‌کند. چه از Databricks، AWS EMR یا Azure Synapse استفاده می‌کنید، تسلط بر دلتا لیک دیگر اختیاری نیست—بلکه یک ضرورت برای مهندسی داده مدرن است.

Share: