Aller au contenu principal

Préparez vos données pour la conformité GDPR

Le Règlement général sur la protection des données (GDPR) et le California Consumer Privacy Act (CCPA) sont des réglementations en matière de confidentialité et de sécurité des données qui exigent des entreprises qu'elles suppriment définitivement et entièrement toutes les informations personnelles identifiables (IPI) collectées sur un client, sur demande explicite de celui-ci. Également connues sous le nom de « droit à l'oubli » (RTBF) ou de « droit à l'effacement des données », les demandes de suppression doivent être exécutées pendant une période spécifiée (par exemple, dans un délai d'un mois calendaire).

Pour implémenter le RTBF sur les données stockées dans Databricks, l'exemple de cet article modélise des datasets pour une entreprise de commerce électronique et montre comment supprimer les données dans les tables sources et propager ces modifications aux tables en aval.

Plan de mise en œuvre du « droit à l'oubli »

Le diagramme suivant illustre comment mettre en œuvre le « droit à l'oubli ».

Schéma qui illustre comment mettre en œuvre la conformité au GDPR.

Suppressions ponctuelles avec Delta Lake

Delta Lake accélère les suppressions ponctuelles dans les grands data lakes avec des transactions ACID, vous permettant de localiser et de supprimer les informations personnelles identifiables (PII) en réponse aux demandes des consommateurs concernant le GDPR ou le CCPA.

Delta Lake conserve l'historique des tables et le rend disponible pour les requêtes à un instant T et les annulations. La fonction VACUUM supprime les fichiers de données qui ne sont plus référencés par une table Delta et qui sont plus anciens qu'un threshold de rétention spécifié, supprimant définitivement les données. Pour en savoir plus sur les valeurs par défaut et les recommandations, consultez Utiliser l'historique de la table.

Assurez-vous que les données sont supprimées lors de l'utilisation de vecteurs de suppression.

Pour les tables avec vecteurs de suppression activés, après la suppression des enregistrements, vous devez également exécuter REORG TABLE ... APPLY (PURGE) pour supprimer définitivement les enregistrements sous-jacents. Ceci inclut les tables Delta Lake, les vues matérialisées et les tables streaming. Voir Appliquer des suppressions logicielles aux fichiers de données.

Supprimer les données dans les sources en amont

Le GDPR et le CCPA s'appliquent à toutes les données, y compris les données dans des sources externes à Delta Lake, telles que Kafka, les fichiers et les bases de données. En plus de supprimer des données dans Databricks, vous devez également penser à supprimer les données dans les sources en amont, telles que les files d'attente et le stockage cloud.

remarque

Avant d’implémenter des workflows de suppression de données, vous devrez peut-être exporter les données du Workspace à des fins de conformité ou de sauvegarde. Consultez Exportation des données du Workspace.

La suppression complète est préférable à l'obfuscation

Vous devez choisir entre supprimer des données et les obfusquer. L’obfuscation peut être mise en œuvre en utilisant la pseudonymisation, le masquage de données, etc. Cependant, l’option la plus sûre est l’effacement complet, car en pratique, l’élimination du risque de réidentification nécessite souvent une suppression complète des données d’informations personnelles identifiables (PII).

Supprimez les données de la couche bronze, puis propagez les suppressions aux couches silver et gold.

Nous vous recommandons de start la conformité GDPR et CCPA par la suppression des données dans la couche bronze en premier, pilotée par un Job planifié qui query une table de demandes de suppression. Après la suppression des données de la couche Bronze, les modifications peuvent être propagées aux couches Silver et Gold.

Entretenez régulièrement les tables pour supprimer les données des fichiers historiques.

Par default, Delta Lake conserve l'historique de la table, y compris les enregistrements supprimés, pendant 30 jours, et le rend disponible pour le time travel et les restaurations. Mais même si les versions précédentes des données sont supprimées, les données sont toujours conservées dans le stockage cloud. Par conséquent, vous devriez régulièrement maintenir les datasets afin de supprimer les versions précédentes des données. La méthode recommandée est l'optimisation prédictive pour les tables gérées par Unity Catalog, qui maintient intelligemment les tables de streaming et les vues matérialisées.

  • Pour les tables gérées par l'optimisation prédictive, les Lakeflow Pipelines maintiennent intelligemment les tables de streaming et les vues matérialisées, en fonction des modèles d'utilisation.
  • Pour les tables sans optimisation prédictive activée, les **LakeFlow Pipelines** effectuent automatiquement des tâches de maintenance dans les 24 heures suivant la mise à jour des tables de streaming et des vues matérialisées.

Si vous n'utilisez pas l'optimisation prédictive ou les LakeFlow Pipelines, vous devriez exécuter une commande VACUUM sur les tables Delta pour supprimer définitivement les versions précédentes des données. By default, cela réduit les capacités de time travel à 7 jours, ce qui est un paramètre configurable, et supprime également les versions historiques des données en question du stockage cloud.

Supprimer les données PII de la couche bronze

En fonction de la conception de votre lakehouse, vous pourriez être en mesure de rompre le link entre les données PII et les données utilisateur non PII. Par exemple, si vous utilisez une clé non naturelle telle que user_id au lieu d'une clé naturelle comme l'e-mail, vous pouvez supprimer les données PII, ce qui laisse les données non PII en place.

Le reste de cet article traite du RTBF en supprimant complètement les enregistrements utilisateur de toutes les tables bronze. Vous pouvez supprimer des données en exécutant une commande DELETE, comme illustré dans le code suivant :

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

Lors de la suppression d’un grand nombre d’enregistrements simultanément, nous recommandons d’utiliser la commande MERGE. Le code ci-dessous part du principe que vous disposez d'une table de contrôle nommée gdpr_control_table qui contient une colonne user_id. Vous insérez un enregistrement dans cette table pour chaque utilisateur qui a demandé le « droit à l’oubli » dans cette table.

La commande MERGE spécifie la condition pour les lignes correspondantes. Dans cet exemple, cela correspond aux enregistrements de target_table avec les enregistrements de gdpr_control_table basés sur le user_id. S'il y a une correspondance (par exemple, un user_id dans le target_table et le gdpr_control_table), la ligne dans le target_table est supprimée. Une fois cette commande MERGE réussie, mettez à jour la table de contrôle pour confirmer que la requête a été traitée.

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
""")

Propager les changements des couches bronze aux couches silver et or

Une fois les données supprimées dans la couche bronze, vous devez propager les modifications aux tables des couches argent et gold.

Vues matérialisées : gestion automatique des suppressions

Les vues matérialisées gèrent automatiquement les suppressions dans les sources. Par conséquent, vous n'avez rien de spécial à faire pour vous assurer qu'une vue matérialisée ne contient pas de données qui ont été supprimées d'une source. Vous devez refresh une vue matérialisée et exécuter la maintenance pour vous assurer que les suppressions sont entièrement traitées.

Une vue matérialisée renvoie toujours le bon résultat, car elle utilise le calcul incrémentiel s'il est moins coûteux qu'une recalculation complète, mais jamais au détriment de l'exactitude. En d'autres termes, la suppression de données d'une source pourrait entraîner une recalculation complète d'une vue matérialisée.

Schéma qui illustre comment gérer automatiquement les suppressions.

Tables de streaming : Supprimer les données et lire la source de streaming à l'aide de skipChangeCommits

Les tables de streaming traitent les données en mode ajout uniquement lorsqu'elles streament depuis des sources de tables Delta. Toute autre opération, telle que la mise à jour ou la suppression d'un enregistrement à partir d'une source de streaming, n'est pas prise en charge et interrompt le stream.

remarque

Pour une implémentation de streaming plus robuste, transmettez plutôt en streaming à partir des flux de modification des tables Delta et gérez les mises à jour et les suppressions dans votre code de traitement. Consultez Gérer les modifications apportées aux tables source Delta Lake.

Diagramme qui illustre comment gérer les suppressions dans sts.

Comme le streaming depuis les tables Delta gère uniquement les nouvelles données, vous devez gérer vous-même les modifications des données. La méthode recommandée consiste à : (1) supprimer les données des tables Delta sources à l'aide de DML, (2) supprimer les données de la table de streaming à l'aide de DML, puis (3) mettre à jour la lecture en streaming pour utiliser skipChangeCommits. Cet indicateur indique que la table de streaming doit ignorer tout ce qui n'est pas une insertion, comme les mises à jour ou les suppressions.

Diagramme qui illustre une méthode de conformité au GDPR qui utilise skipChangeCommits.

Vous pouvez également (1) supprimer les données de la source, puis (2) refresh complètement la table de streaming. Lorsque vous effectuez un refresh complet d’une table de streaming, cela efface l’état de streaming de la table et retraitera toutes les données. Toute source de données en amont dont la période de conservation est dépassée (par exemple, un sujet Kafka qui supprime les données après 7 jours) ne sera plus traitée, ce qui pourrait entraîner une perte de données. Nous recommandons cette option pour les tables de streaming uniquement dans le scénario où les données historiques sont disponibles et que leur traitement à nouveau ne sera pas coûteux.

Schéma illustrant une méthode de conformité GDPR qui effectue un refresh complet sur le st.

Exemple : conformité GDPR et CCPA pour une entreprise de commerce électronique

Le diagramme suivant présente une architecture en médaillon pour une entreprise de commerce électronique où la conformité GDPR & CCPA doit être mise en œuvre. Même si les données d'un utilisateur sont supprimées, vous pourriez vouloir compter leurs activités dans les agrégations en aval.

Diagramme illustrant un exemple de conformité GDPR et CCPA pour une entreprise d'e-commerce.

  • Tables sources

    • source_users - Une table source streaming d'utilisateurs (créée ici, pour l'exemple). Les environnements de production utilisent généralement Kafka, Kinesis ou des plateformes de streaming similaires.
    • source_clicks - Une table source de streaming de clics (créée ici, pour l'exemple). Les environnements de production utilisent généralement Kafka, Kinesis ou des plateformes de streaming similaires.
  • Table de contrôle

    • gdpr_requests - Table de contrôle contenant les ID utilisateur soumis au « droit à l'oubli ». Lorsqu'un utilisateur demande à être supprimé, ajoutez-le ici.
  • Couche Bronze

    • users_bronze – Dimensions utilisateur. Contient des informations personnelles identifiables (par exemple, une adresse e-mail).
    • clicks_bronze - Cliquez sur événements. Contient des PII (par exemple, adresse IP).
  • Couche Silver

    • clicks_silver - Données de clics nettoyées et standardisées.
    • users_silver - Données utilisateur nettoyées et standardisées.
    • user_clicks_silver - Joint clicks_silver (streaming) avec un instantané de users_silver.
  • Couche Gold

    • user_behavior_gold – Métriques agrégées du comportement de l'utilisateur.
    • marketing_insights_gold - Segment d'utilisateurs pour les insights du marché.

Étape 1 : Remplir les tables avec un échantillon de données

Le code suivant crée ces deux tables pour cet exemple et les remplit avec des exemples de données :

  • source_users contient des données dimensionnelles sur les utilisateurs. Cette table contient une colonne PII appelée email.
  • source_clicks contient des données d'événement concernant les activités réalisées par les utilisateurs. Elle contient une colonne PII nommée ip_address.
Python
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")

Étape 2 : créer un pipeline qui traite les données PII.

Le code suivant crée les couches bronze, argent et gold de l'architecture en médaillon présentée ci-dessus.

Python
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')
)
)

Étape 3 : Supprimer les données dans les tables sources

Au cours de cette étape, vous supprimez les données de toutes les tables où des PII sont trouvées. La fonction suivante supprime toutes les instances des PII d'un utilisateur des tables contenant des PII.

Python
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}
""")

Étape 4 : Ajouter skipChangeCommits aux définitions des tables de streaming concernées

Dans cette étape, vous devez dire à la pipeline d'ignorer les lignes non ajoutées. Ajoutez l’option skipChangeCommits aux méthodes suivantes. Vous n'avez pas à mettre à jour les définitions des vues matérialisées, car elles gèrent automatiquement les mises à jour et les suppressions.

  • users_bronze
  • users_silver
  • clicks_bronze
  • clicks_silver
  • user_clicks_silver

Le code suivant montre comment mettre à jour la méthode users_bronze :

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

Lorsque vous exécutez à nouveau le pipeline, la mise à jour réussit.