Production pipelines with workflows

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Disha Mukherjee

Lead Data Engineer

ทำไมต้องใช้ Delta Lake?

คลังข้อมูลดิจิทัลที่ปลอดภัย มีตารางข้อมูลเป็นระเบียบและโล่ป้องกัน สไตล์แบน ทันสมัย

 

$$

  • ACID transactions → ย้อนกลับการเขียนที่ล้มเหลว
  • Schema enforcement → บล็อกประเภทข้อมูลที่ไม่ตรงกัน
  • Versioning → ดึงข้อมูลสถานะก่อนหน้าได้ทุกเวอร์ชัน
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การเขียนลง Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • ตาราง Delta ใหม่จะปรากฏใน Unity Catalog
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Notebook tasks

task1_ingest: โหลดและทำความสะอาดข้อมูล

  • โหลด CSV → ทำความสะอาด → เขียนเป็นตาราง Delta
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Notebook tasks

task2_metrics: รายได้แยกตามหมวดหมู่

  • อ่านตารางที่ทำความสะอาดแล้ว → คำนวณ metrics → เขียนลงตารางใหม่
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Notebook tasks

task3_customers: จัดอันดับตามยอดใช้จ่าย

  • อ่านตารางที่ทำความสะอาดแล้ว → จัดอันดับลูกค้า → บันทึกลงตารางใหม่
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การสร้าง job

UI ของ Jobs and Pipelines แสดง job DAG สามงานที่เสร็จสมบูรณ์ โดยมีลูกศรลำดับการทำงานจาก task1_ingest ไป task2_metrics และ task3_customers

  • รันด้วยตนเอง กำหนดเวลา หรือตั้ง trigger ได้
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การรัน job

Job run DAG ที่ task1_ingest และ task2_metrics สำเร็จ (สีเขียว) และ task3_customers ล้มเหลว (สีแดง)

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การรัน job

ตรวจสอบข้อผิดพลาดใน task สาม

  • ข้อผิดพลาด: customer-id ควรเป็น Customer_ID
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การรัน job

มุมมอง graph แบบสำเร็จ

มุมมอง timeline

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Lakeflow คืออะไร?

$$

$$

เปรียบเทียบ: Imperative (Jobs) | Declarative (Lakeflow)

 

  • Jobs → จัดการแต่ละขั้นตอนเอง
  • Lakeflow → ระบุเพียงว่าตารางควรมีข้อมูลอะไร
  • Databricks จัดการลำดับ การ retry และ compute ให้
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

รูปแบบ @dlt.table

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

การรัน pipeline

Lakeflow pipeline DAG แสดง transactions_bronze 100K แถว ต่อด้วย transactions_silver 33K แถว และ category_revenue_gold 6 แถว ทั้งหมดเป็น materialized views สีเขียว

  • Bronze (100K แถว) → Silver (33K แถว) → Gold (6 แถว)
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

ควรใช้ Notebooks, Jobs หรือ Lakeflow?

$$

ชั้น: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → สำรวจและสร้าง prototype
  • Databricks Jobs → pipeline หลายขั้นตอน กำหนดเวลาได้
  • Lakeflow → จัดการครบวงจร แบบ declarative
การแปลงข้อมูลด้วย Spark SQL ใน Databricks

มาฝึกกันเถอะ!

การแปลงข้อมูลด้วย Spark SQL ใน Databricks

Preparing Video For Download...