Tutorial: Buat pipeline pertama Anda menggunakan Editor Lakeflow Pipelines

Anda membuat alur Lakeflow baru untuk orkestrasi data dengan Auto Loader, lalu memperluas alur sampel dengan membersihkan data dan membuat kueri untuk menemukan 100 pengguna teratas.

Dalam tutorial ini, Anda mempelajari cara menggunakan Editor Alur Lakeflow untuk:

  • Buat alur baru dengan struktur folder default dan mulai dengan sekumpulan file sampel.
  • Tentukan batasan kualitas data menggunakan harapan.
  • Gunakan fitur editor untuk memperluas alur dengan transformasi baru untuk melakukan analisis pada data Anda.

Persyaratan

Sebelum memulai tutorial ini, Anda harus:

  • Masuk ke ruang kerja Azure Databricks.
  • Mengaktifkan Unity Catalog untuk ruang kerja Anda.
  • Memiliki izin untuk membuat sumber daya komputasi atau akses ke sumber daya komputasi.
  • Memiliki izin untuk membuat skema baru dalam katalog. Izin yang diperlukan adalah ALL PRIVILEGES atau USE CATALOG dan CREATE SCHEMA.
  • Untuk melihat seluruh hak istimewa yang diperlukan untuk membuat, menjalankan, memperbarui, dan melihat pipeline serta outputnya, lihat Mengelola identitas, izin, dan hak istimewa untuk pipeline.

Langkah 1: Membuat alur

Dalam langkah ini, Anda membuat alur menggunakan struktur folder default dan sampel kode. Sampel kode mereferensikan users tabel di wanderbricks sumber data sampel.

  1. Di ruang kerja Azure Databricks Anda, klik Plus icon.Baru, kemudian Pipeline icon.ETL. Ini membuka editor alur dengan nama alur default seperti New Pipeline <date> <time>.

  2. (Opsional) Pilih nama dan masukkan nama deskriptif untuk alur.

  3. (Opsional) Di sebelah kanan nama, klik katalog dan skema untuk mengatur default yang berbeda.

  4. (Opsional) Dalam file sumber my_transformation yang dibuat untuk Anda, pilih Python atau SQL dari daftar drop-down bahasa untuk mengatur bahasa file.

  5. Klik ikon Kode.Gunakan kode sampel.

    Kode sampel dalam bahasa yang Anda pilih muncul di my_transformation file sumber di transformations folder . Himpunan data output belum dibuat, dan grafik Alur di sisi kanan layar kosong.

  6. Untuk menjalankan kode alur (kode dalam transformations folder), klik Jalankan alur di bagian kanan atas layar.

    Setelah proses selesai, bagian bawah ruang kerja memperlihatkan dua tabel baru yang dibuat, sample_users_<date_time> dan sample_aggregation_<date_time>. Grafik Alur di sisi kanan ruang kerja sekarang menunjukkan dua tabel, termasuk yang sample_users merupakan sumber untuk sample_aggregation. Catat nama tabel lengkap sample_users_<date_time> ; Anda mereferensikannya di langkah berikutnya.

Langkah 2: Menerapkan pemeriksaan kualitas data

Dalam langkah ini, Anda menambahkan pemeriksaan kualitas data ke sample_users tabel. Anda menggunakan ekspektasi alur untuk membatasi data. Dalam hal ini, Anda menghapus catatan pengguna apa pun yang tidak memiliki alamat email yang valid, dan menghasilkan tabel yang dibersihkan sebagai users_cleaned.

  1. Pada browser aset pipeline di sebelah kiri, klik ikon Plus, lalu pilih Transformasi.

  2. Dalam dialog Buat file transformasi baru , buat pilihan berikut:

    • Pilih Python atau SQL untuk Language. Ini tidak harus cocok dengan pilihan Anda sebelumnya.
    • Beri nama file. Dalam hal ini, pilih users_cleaned.
    • Untuk Jalur tujuan, biarkan default.
    • Untuk Jenis himpunan data, biarkan sebagai Tidak Ada yang dipilih atau pilih Tampilan materialisasi. Jika Anda memilih Tampilan materialisasi, itu menghasilkan kode sampel untuk Anda.
  3. Klik Buat untuk membuat file kode transformasi.

  4. Dalam file kode baru Anda, edit kode agar sesuai dengan yang berikut (gunakan SQL atau Python, berdasarkan pilihan Anda di layar sebelumnya). Ganti sample_users_<date_time> dengan nama lengkap tabel Anda sample_users dari bagian sebelumnya.

    SQL

    -- Drop all rows that do not have an email address
    
    CREATE MATERIALIZED VIEW users_cleaned
    (
      CONSTRAINT non_null_email EXPECT (email IS NOT NULL) ON VIOLATION DROP ROW
    ) AS
    SELECT *
    FROM sample_users_<date_time>;
    

    Python

    from pyspark import pipelines as dp
    
    # Drop all rows that do not have an email address
    
    @dp.materialized_view
    @dp.expect_or_drop("no null emails", "email IS NOT NULL")
    def users_cleaned():
        return (
            spark.read.table("sample_users_<date_time>")
        )
    
  5. Klik Jalankan alur untuk memperbarui alur. Sekarang harus memiliki tiga tabel.

Langkah 3: Menganalisis pengguna teratas

Selanjutnya dapatkan 100 pengguna teratas dengan jumlah pemesanan yang telah mereka buat. Gabungkan wanderbricks.bookings tabel ke tampilan materialisasi users_cleaned .

  1. Di browser aset pipeline di sebelah kiri, klik ikon Plus, lalu pilih opsi Transformasi.

  2. Dalam dialog Buat file transformasi baru , buat pilihan berikut:

    • Pilih Python atau SQL untuk Language. Ini tidak harus cocok dengan pilihan Anda sebelumnya.
    • Beri nama file. Dalam hal ini, pilih users_and_bookings.
    • Untuk Jalur tujuan, biarkan default.
    • Untuk Jenis himpunan data, biarkan sebagai Tidak Ada yang dipilih.
  3. Klik Buat untuk membuat file kode transformasi.

  4. Dalam file kode baru Anda, edit kode agar sesuai dengan yang berikut (gunakan SQL atau Python, berdasarkan pilihan Anda di layar sebelumnya).

    SQL

    -- Get the top 100 users by number of bookings
    
    CREATE OR REFRESH MATERIALIZED VIEW users_and_bookings AS
    SELECT u.name AS name, COUNT(b.booking_id) AS booking_count
    FROM users_cleaned u
    JOIN samples.wanderbricks.bookings b ON u.user_id = b.user_id
    GROUP BY u.name
    ORDER BY booking_count DESC
    LIMIT 100;
    

    Python

    from pyspark import pipelines as dp
    from pyspark.sql.functions import col, count, desc
    
    # Get the top 100 users by number of bookings
    
    @dp.materialized_view
    def users_and_bookings():
        return (
            spark.read.table("users_cleaned")
            .join(spark.read.table("samples.wanderbricks.bookings"), "user_id")
            .groupBy(col("name"))
            .agg(count("booking_id").alias("booking_count"))
            .orderBy(desc("booking_count"))
            .limit(100)
        )
    
  5. Klik Jalankan alur untuk memperbarui himpunan data. Saat proses selesai, Anda dapat melihat di Grafik Alur bahwa ada empat tabel, termasuk tabel baru users_and_bookings .

    Grafik alur memperlihatkan empat tabel dalam alur

Sumber daya tambahan

Sekarang setelah Anda mempelajari cara menggunakan beberapa fitur editor alur Lakeflow dan membuat alur, berikut adalah beberapa fitur lain untuk mempelajari selengkapnya tentang: