Pipeline produksi dengan workflows

Transformasi Data dengan Spark SQL di Databricks

Disha Mukherjee

Lead Data Engineer

Mengapa Delta Lake?

Sebuah brankas digital bercahaya yang aman dengan tabel data teratur dan perisai pelindung, gaya modern flat

 

$$

  • Transaksi ACID → membatalkan tulis yang gagal
  • Penegakan skema → memblokir tipe yang tidak cocok
  • Versioning → kueri keadaan sebelumnya mana pun
Transformasi Data dengan Spark SQL di Databricks

Menulis ke Delta

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

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

$$

$$

  • Tabel Delta baru muncul di Unity Catalog
Transformasi Data dengan Spark SQL di Databricks

Tugas notebook

task1_ingest: muat dan bersihkan

  • Memuat CSV → menerapkan pembersihan → menulis tabel Delta
Transformasi Data dengan Spark SQL di Databricks

Tugas notebook

task2_metrics: pendapatan per kategori

  • Membaca tabel bersih → menghitung metrik → menulis ke tabel baru
Transformasi Data dengan Spark SQL di Databricks

Tugas notebook

task3_customers: peringkat berdasarkan belanja

  • Membaca tabel bersih → memberi peringkat pelanggan → menyimpan ke tabel baru
Transformasi Data dengan Spark SQL di Databricks

Membuat job

UI Jobs and Pipelines menampilkan DAG pekerjaan tiga tugas yang selesai dengan panah ketergantungan dari task1_ingest ke task2_metrics ke task3_customers

  • Dapat dijalankan manual, dijadwalkan, dan diatur pemicunya
Transformasi Data dengan Spark SQL di Databricks

Menjalankan job

DAG run job dengan task1_ingest dan task2_metrics berhasil (hijau) dan task3_customers gagal (merah)

Transformasi Data dengan Spark SQL di Databricks

Menjalankan job

Memeriksa error di tugas ketiga

  • Error: customer-id seharusnya Customer_ID
Transformasi Data dengan Spark SQL di Databricks

Menjalankan job

Tampilan grafik dengan sukses

Tampilan linimasa

Transformasi Data dengan Spark SQL di Databricks

Apa itu Lakeflow?

$$

$$

perbandingan: Imperatif (Jobs) | Deklaratif (Lakeflow)

 

  • Jobs → kita mengelola tiap langkah
  • Lakeflow → nyatakan apa isi tabel
  • Databricks menangani urutan, percobaan ulang, dan komputasi
Transformasi Data dengan Spark SQL di Databricks

Pola @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(...)
Transformasi Data dengan Spark SQL di Databricks

Menjalankan pipeline

DAG pipeline Lakeflow menampilkan transactions_bronze 100K baris lalu transactions_silver 33K baris lalu category_revenue_gold 6 baris semuanya hijau materialized views

  • Bronze (100K baris) → Silver (33K baris) → Gold (6 baris)
Transformasi Data dengan Spark SQL di Databricks

Notebooks, Jobs, atau Lakeflow?

$$

lapisan: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → eksplorasi, prototipe
  • Databricks Jobs → pipeline multi-langkah terjadwal
  • Lakeflow → dikelola penuh, deklaratif
Transformasi Data dengan Spark SQL di Databricks

Ayo berlatih!

Transformasi Data dengan Spark SQL di Databricks

Preparing Video For Download...