Menyiapkan data Anda untuk kepatuhan GDPR

Peraturan Perlindungan Data Umum (GDPR) dan Undang-Undang Privasi Konsumen California (CCPA) adalah peraturan privasi dan keamanan data yang mengharuskan perusahaan untuk menghapus semua informasi identitas pribadi (PII) secara permanen dan sepenuhnya yang dikumpulkan tentang pelanggan atas permintaan eksplisit mereka. Juga dikenal sebagai "hak untuk dilupakan" (RTBF) atau "hak untuk penghapusan data", permintaan penghapusan harus dijalankan selama periode tertentu (misalnya, dalam satu bulan kalender).

Untuk menerapkan RTBF pada data yang disimpan di Azure Databricks, contoh dalam artikel ini memodelkan himpunan data untuk perusahaan e-niaga dan menunjukkan cara menghapus data dalam tabel sumber dan menyebarluaskan perubahan tersebut ke tabel hilir.

Cetak biru untuk menerapkan "hak untuk dilupakan"

Diagram berikut menggambarkan cara mengimplementasikan "hak untuk dilupakan."

Diagram yang menggambarkan cara menerapkan kepatuhan GDPR.

Penghapusan titik dengan Delta Lake

Delta Lake mempercepat penghapusan data pada titik tertentu di data lake berukuran besar dengan transaksi ACID, sehingga Anda dapat menemukan dan menghapus informasi identitas pribadi (PII) sebagai tanggapan atas permintaan konsumen terkait GDPR atau CCPA.

Delta Lake mempertahankan riwayat tabel dan membuatnya tersedia untuk kueri dan pemulihan ke titik waktu tertentu. Fungsi VACUUM menghapus file data yang tidak lagi direferensikan oleh tabel Delta dan lebih lama dari ambang retensi yang ditentukan, menghapus data secara permanen. Untuk mempelajari selengkapnya tentang default dan rekomendasi, lihat Bekerja dengan riwayat tabel.

Pastikan data dihapus saat menggunakan vektor penghapusan

Untuk tabel dengan vektor penghapusan diaktifkan, setelah menghapus rekaman, Anda juga harus menjalankan REORG TABLE ... APPLY (PURGE) untuk menghapus rekaman yang mendasarinya secara permanen. Ini termasuk tabel Delta Lake, tampilan materialisasi, dan tabel streaming. Lihat Menerapkan penghapusan sementara ke file data.

Menghapus data di sumber hulu

GDPR dan CCPA berlaku untuk semua data, termasuk data dalam sumber di luar Delta Lake, seperti Kafka, file, dan database. Selain menghapus data di Databricks, Anda juga harus ingat untuk menghapus data di sumber hulu, seperti antrean dan penyimpanan cloud.

Note

Sebelum menerapkan alur kerja penghapusan data, Anda mungkin perlu mengekspor data ruang kerja untuk tujuan kepatuhan atau pencadangan. Lihat Mengekspor data ruang kerja.

Penghapusan lengkap lebih diutamakan daripada pengaburan

Anda harus memilih antara menghapus data dan mengaburkannya. Obfuscation dapat diimplementasikan menggunakan pseudonimisasi, masking data, dll. Namun, opsi paling aman adalah penghapusan lengkap karena, dalam praktiknya, menghilangkan risiko reidentifikasi sering memerlukan penghapusan lengkap data PII.

Hapus data dalam lapisan perunggu, lalu sebarkan penghapusan ke lapisan perak dan emas

Kami menyarankan agar Anda memulai kepatuhan GDPR dan CCPA dengan menghapus data di lapisan perunggu terlebih dahulu, didorong oleh pekerjaan terjadwal yang meminta tabel permintaan penghapusan. Setelah data dihapus dari lapisan perunggu, perubahan dapat disebarkan ke lapisan perak dan emas.

Memelihara tabel secara teratur untuk menghapus data dari file historis

Secara default, Delta Lake mempertahankan riwayat tabel, termasuk rekaman yang dihapus, selama 30 hari, dan membuatnya tersedia untuk perjalanan waktu dan pemutaran kembali. Tetapi bahkan jika versi data sebelumnya dihapus, data masih disimpan di penyimpanan cloud. Oleh karena itu, Anda harus secara teratur memelihara himpunan data untuk menghapus versi data sebelumnya. Cara yang disarankan adalah Pengoptimalan Prediktif untuk tabel terkelola Unity Catalog, yang secara cerdas mempertahankan tabel streaming dan tampilan materialisasi.

  • Untuk tabel yang dikelola oleh pengoptimalan prediktif, alur Lakeflow secara cerdas mempertahankan tabel streaming dan tampilan materialisasi, berdasarkan pola penggunaan.
  • Untuk tabel tanpa pengoptimalan prediktif diaktifkan, alur Lakeflow secara otomatis melakukan tugas pemeliharaan dalam waktu 24 jam setelah tabel streaming dan tampilan materialisasi diperbarui.

Jika Anda tidak menggunakan pengoptimalan prediktif atau alur Lakeflow, Anda harus menjalankan VACUUM perintah pada tabel Delta untuk menghapus versi data sebelumnya secara permanen. Secara default, ini mengurangi kemampuan perjalanan waktu menjadi 7 hari, yang merupakan pengaturan yang dapat dikonfigurasi, dan menghapus versi historis data yang dimaksud dari penyimpanan cloud juga.

Menghapus data PII dari lapisan perunggu

Tergantung pada desain lakehouse Anda, Anda mungkin dapat memutus tautan antara data pengguna PII dan non-PII. Misalnya, jika Anda menggunakan kunci non-alami seperti user_id alih-alih kunci alami seperti email, Anda dapat menghapus data PII, yang meninggalkan data non-PII.

Sisa artikel ini menangani RTBF dengan menghapus catatan pengguna sepenuhnya dari semua tabel perunggu. Anda dapat menghapus data dengan menjalankan perintah DELETE, seperti yang diperlihatkan dalam kode berikut:

spark.sql("DELETE FROM bronze.users WHERE user_id = 5")

Saat menghapus sejumlah besar rekaman bersama-sama pada satu waktu, sebaiknya gunakan perintah MERGE. Kode di bawah ini mengasumsikan bahwa Anda memiliki tabel kontrol yang disebut gdpr_control_table yang berisi kolom user_id. Anda menyisipkan rekaman ke dalam tabel ini untuk setiap pengguna yang telah meminta "hak untuk dilupakan" ke dalam tabel ini.

Perintah MERGE menentukan kondisi untuk baris yang cocok. Dalam contoh ini, mencocokkan rekaman dari target_table dengan rekaman di gdpr_control_table berdasarkan user_id. Jika ada kecocokan (misalnya, user_id di target_table dan gdpr_control_table), baris di target_table dihapus. Setelah perintah MERGE ini berhasil, perbarui tabel kontrol untuk mengonfirmasi bahwa permintaan telah diproses.

spark.sql("""
  MERGE INTO target
  USING (
    SELECT user_id
    FROM gdpr_control_table
  ) AS source
  ON target.user_id = source.user_id
  WHEN MATCHED THEN DELETE
""")

Menyebarluaskan perubahan dari lapisan perunggu ke lapisan perak dan emas

Setelah data dihapus di lapisan perunggu, Anda harus menyebarluaskan perubahan pada tabel di lapisan perak dan emas.

Tampilan materialisasi: Menangani penghapusan secara otomatis

Tampilan materialisasi menangani penghapusan pada sumber secara otomatis. Oleh karena itu, Anda tidak perlu melakukan sesuatu yang istimewa untuk memastikan bahwa tampilan materialisasi tidak berisi data yang telah dihapus dari sumber. Anda harus memperbarui tampilan materialisasi dan menjalankan pemeliharaan untuk memastikan bahwa penghapusan diproses sepenuhnya.

Tampilan termaterialisasi selalu mengembalikan hasil yang benar karena menggunakan komputasi inkremental jika lebih murah daripada komputasi ulang penuh, namun tidak pernah dengan mengorbankan kebenaran. Dengan kata lain, menghapus data dari sumber dapat menyebabkan tampilan materialisasi melakukan komputasi ulang sepenuhnya.

Diagram yang menggambarkan cara menangani penghapusan secara otomatis.

Tabel streaming: Menghapus data dan membaca sumber streaming menggunakan skipChangeCommits

Tabel streaming memproses data khusus tambahan saat mereka melakukan streaming dari sumber tabel Delta. Operasi lain apa pun, seperti memperbarui atau menghapus rekaman dari sumber streaming, tidak didukung dan merusak aliran.

Note

Untuk implementasi streaming yang lebih kuat, streaming dari umpan perubahan tabel Delta sebagai gantinya, dan tangani pembaruan dan penghapusan dalam kode pemrosesan Anda. Lihat Menangani perubahan pada tabel Delta Lake sumber.

Diagram yang menggambarkan cara menangani penghapusan di sts.

Karena streaming dari tabel Delta hanya menangani data baru, Anda harus menangani perubahan pada data sendiri. Metode yang disarankan adalah untuk: (1) menghapus data dalam tabel Delta sumber menggunakan DML, (2) menghapus data dari tabel streaming menggunakan DML, lalu (3) memperbarui streaming yang dibaca untuk digunakan skipChangeCommits. Penanda ini menunjukkan bahwa tabel streaming harus mengabaikan segala sesuatu selain penyisipan, seperti pembaruan atau penghapusan.

Diagram yang mengilustrasikan metode kepatuhan GDPR yang menggunakan skipChangeCommits.

Atau, Anda dapat (1) menghapus data dari sumbernya, lalu (2) sepenuhnya me-refresh tabel streaming. Ketika Anda melakukan segar ulang penuh pada tabel streaming, itu akan menghapus status streaming tabel dan memproses kembali seluruh data. Setiap sumber data hulu yang berada di luar periode retensinya (misalnya, topik Kafka yang menua setelah 7 hari) tidak akan diproses lagi, yang dapat menyebabkan kehilangan data. Kami merekomendasikan opsi ini untuk tabel streaming hanya dalam skenario di mana data historis tersedia dan memprosesnya lagi tidak akan mahal.

Diagram yang menggambarkan metode kepatuhan GDPR yang melakukan penyegaran penuh pada st.

Contoh: Kepatuhan GDPR dan CCPA untuk perusahaan e-niaga

Diagram berikut menunjukkan arsitektur medali untuk perusahaan e-niaga, di mana kepatuhan terhadap GDPR & CCPA perlu diterapkan. Meskipun data pengguna dihapus, Anda mungkin ingin menghitung aktivitas mereka dalam agregasi hilir.

Diagram yang menggambarkan contoh kepatuhan GDPR dan CCPA untuk perusahaan e-niaga.

  • Tabel sumber
    • source_users - Tabel sumber streaming pengguna (dibuat di sini, misalnya). Lingkungan produksi biasanya menggunakan Kafka, Kinesis, atau platform streaming serupa.
    • source_clicks - Tabel klik sumber streaming (dibuat di sini, misalnya). Lingkungan produksi biasanya menggunakan Kafka, Kinesis, atau platform streaming serupa.
  • Tabel kontrol
    • gdpr_requests - Tabel kontrol yang berisi ID pengguna tunduk pada "hak untuk dilupakan." Saat pengguna meminta untuk dihapus, tambahkan di sini.
  • lapisan Perunggu
    • users_bronze - Dimensi pengguna. Berisi PII (misalnya, alamat email).
    • clicks_bronze - Klik peristiwa. Berisi PII (misalnya, alamat IP).
  • lapisan perak
    • clicks_silver - Data klik yang dibersihkan dan distandarkan.
    • users_silver - Data pengguna yang dibersihkan dan distandarkan.
    • user_clicks_silver- Menggabungkan clicks_silver (streaming) dengan rekaman jepret dari users_silver.
  • lapisan Emas
    • user_behavior_gold - Metrik perilaku pengguna agregat.
    • marketing_insights_gold - Segmen pengguna untuk wawasan pasar.

Langkah 1: Mengisi tabel dengan data sampel

Kode berikut membuat dua tabel ini untuk contoh ini dan mengisinya dengan data sampel:

  • source_users berisi data dimensi tentang pengguna. Tabel ini berisi kolom PII yang disebut email.
  • source_clicks berisi data peristiwa tentang aktivitas yang dilakukan oleh pengguna. Ini berisi kolom PII yang disebut ip_address.
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, MapType, DateType

catalog = "users"
schema = "name"

# Create table containing sample users
users_schema = StructType([
   StructField('user_id', IntegerType(), False),
   StructField('username', StringType(), True),
   StructField('email', StringType(), True),
   StructField('registration_date', StringType(), True),
   StructField('user_preferences', MapType(StringType(), StringType()), True)
])

users_data = [
   (1, 'alice', 'alice@example.com', '2021-01-01', {'theme': 'dark', 'language': 'en'}),
   (2, 'bob', 'bob@example.com', '2021-02-15', {'theme': 'light', 'language': 'fr'}),
   (3, 'charlie', 'charlie@example.com', '2021-03-10', {'theme': 'dark', 'language': 'es'}),
   (4, 'david', 'david@example.com', '2021-04-20', {'theme': 'light', 'language': 'de'}),
   (5, 'eve', 'eve@example.com', '2021-05-25', {'theme': 'dark', 'language': 'it'})
]

users_df = spark.createDataFrame(users_data, schema=users_schema)
users_df.write.mode("overwrite").saveAsTable(f"{catalog}.{schema}.source_users")

# Create table containing clickstream (i.e. user activities)
from pyspark.sql.types import TimestampType

clicks_schema = StructType([
   StructField('click_id', IntegerType(), False),
   StructField('user_id', IntegerType(), True),
   StructField('url_clicked', StringType(), True),
   StructField('click_timestamp', StringType(), True),
   StructField('device_type', StringType(), True),
   StructField('ip_address', StringType(), True)
])

clicks_data = [
   (1001, 1, 'https://example.com/home', '2021-06-01T12:00:00', 'mobile', '192.168.1.1'),
   (1002, 1, 'https://example.com/about', '2021-06-01T12:05:00', 'desktop', '192.168.1.1'),
   (1003, 2, 'https://example.com/contact', '2021-06-02T14:00:00', 'tablet', '192.168.1.2'),
   (1004, 3, 'https://example.com/products', '2021-06-03T16:30:00', 'mobile', '192.168.1.3'),
   (1005, 4, 'https://example.com/services', '2021-06-04T10:15:00', 'desktop', '192.168.1.4'),
   (1006, 5, 'https://example.com/blog', '2021-06-05T09:45:00', 'tablet', '192.168.1.5')
]

clicks_df = spark.createDataFrame(clicks_data, schema=clicks_schema)
clicks_df.write.format("delta").mode("overwrite").saveAsTable(f"{catalog}.{schema}.source_clicks")

Langkah 2: Membuat alur yang memproses data PII

Kode berikut membuat lapisan perunggu, perak, dan emas dari arsitektur medali yang ditunjukkan di atas.

from pyspark import pipelines as dp
from pyspark.sql.functions import col, concat_ws, count, countDistinct, avg, when, expr

catalog = "users"
schema = "name"

# ----------------------------
# Bronze Layer - Raw Data Ingestion
# ----------------------------

@dp.table(
   name=f"{catalog}.{schema}.users_bronze",
   comment='Raw users data loaded from source'
)
def users_bronze():
   return (
     spark.readStream.table(f"{catalog}.{schema}.source_users")
   )

@dp.table(
   name=f"{catalog}.{schema}.clicks_bronze",
   comment='Raw clicks data loaded from source'
)
def clicks_bronze():
   return (
       spark.readStream.table(f"{catalog}.{schema}.source_clicks")
   )

# ----------------------------
# Silver Layer - Data Cleaning and Enrichment
# ----------------------------

@dp.create_streaming_table(
   name=f"{catalog}.{schema}.users_silver",
   comment='Cleaned and standardized users data'
)

@dp.view
@dp.expect_or_drop('valid_email', "email IS NOT NULL")
def users_bronze_view():
   return (
       spark.readStream
           .table(f"{catalog}.{schema}.users_bronze")
           .withColumn('registration_date', col('registration_date').cast('timestamp'))
           .dropDuplicates(['user_id', 'registration_date'])
           .select('user_id', 'username', 'email', 'registration_date', 'user_preferences')
   )

@dp.create_auto_cdc_flow(
   target=f"{catalog}.{schema}.users_silver",
   source="users_bronze_view",
   keys=["user_id"],
   sequence_by="registration_date",
)

@dp.table(
   name=f"{catalog}.{schema}.clicks_silver",
   comment='Cleaned and standardized clicks data'
)
@dp.expect_or_drop('valid_click_timestamp', "click_timestamp IS NOT NULL")
def clicks_silver():
   return (
       spark.readStream
           .table(f"{catalog}.{schema}.clicks_bronze")
           .withColumn('click_timestamp', col('click_timestamp').cast('timestamp'))
           .withWatermark('click_timestamp', '10 minutes')
           .dropDuplicates(['click_id'])
           .select('click_id', 'user_id', 'url_clicked', 'click_timestamp', 'device_type', 'ip_address')
   )

@dp.table(
   name=f"{catalog}.{schema}.user_clicks_silver",
   comment='Joined users and clicks data on user_id'
)
def user_clicks_silver():
   # Read users_silver as a static DataFrame - each refresh
   # will use a snapshot of the users_silver table.
   users = spark.read.table(f"{catalog}.{schema}.users_silver")

   # Read clicks_silver as a streaming DataFrame.
   clicks = spark.readStream \
       .table('clicks_silver')

   # Perform the join - join of a static dataset with a
   # streaming dataset creates a streaming table.
   joined_df = clicks.join(users, on='user_id', how='inner')

   return joined_df

# ----------------------------
# Gold Layer - Aggregated and Business-Level Data
# ----------------------------

@dp.materialized_view(
   name=f"{catalog}.{schema}.user_behavior_gold",
   comment='Aggregated user behavior metrics'
)
def user_behavior_gold():
   df = spark.read.table(f"{catalog}.{schema}.user_clicks_silver")
   return (
       df.groupBy('user_id')
         .agg(
             count('click_id').alias('total_clicks'),
             countDistinct('url_clicked').alias('unique_urls')
         )
   )

@dp.materialized_view(
   name=f"{catalog}.{schema}.marketing_insights_gold",
   comment='User segments for marketing insights'
)
def marketing_insights_gold():
   df = spark.read.table(f"{catalog}.{schema}.user_behavior_gold")
   return (
       df.withColumn(
           'engagement_segment',
           when(col('total_clicks') >= 100, 'High Engagement')
           .when((col('total_clicks') >= 50) & (col('total_clicks') < 100), 'Medium Engagement')
           .otherwise('Low Engagement')
       )
   )

Langkah 3: Menghapus data dalam tabel sumber

Dalam langkah ini, Anda menghapus data di semua tabel tempat PII ditemukan. Fungsi berikut menghapus semua instans PII pengguna dari tabel dengan PII.

catalog = "users"
schema = "name"

def apply_gdpr_delete(user_id):
 tables_with_pii = ["clicks_bronze", "users_bronze", "clicks_silver", "users_silver", "user_clicks_silver"]

 for table in tables_with_pii:
   print(f"Deleting user_id {user_id} from table {table}")
   spark.sql(f"""
     DELETE FROM {catalog}.{schema}.{table}
     WHERE user_id = {user_id}
   """)

Langkah 4: Tambahkan skipChangeCommits ke definisi tabel streaming yang terpengaruh

Dalam langkah ini, Anda harus memberi tahu alur untuk melewati baris yang tidak ditambahkan. Tambahkan opsi skipChangeCommits ke metode berikut. Anda tidak perlu memperbarui definisi tampilan materialisasi karena secara otomatis menangani pembaruan dan penghapusan.

  • users_bronze
  • users_silver
  • clicks_bronze
  • clicks_silver
  • user_clicks_silver

Kode berikut menunjukkan cara memperbarui metode users_bronze:

def users_bronze():
   return (
     spark.readStream.option('skipChangeCommits', 'true').table(f"{catalog}.{schema}.source_users")
   )

Saat Anda menjalankan pipeline lagi, pembaruan berhasil dilakukan.