Aller au contenu principal

Évolution des schémas dans Databricks

L' évolution des schémas fait référence à la capacité d'un système à s'adapter aux changements dans la structure des données au fil du temps. Ces changements sont fréquents lorsque vous travaillez avec des données semi-structurées, des event Stream ou des sources tierces où de nouveaux champs sont ajoutés, les types de données changent ou les structures imbriquées évoluent.

Les changements courants incluent :

  • Nouvelles colonnes : champs supplémentaires non définis précédemment, parfois avec une valeur de remplissage personnalisée.
  • Renommage de colonne : Modification d'un nom de colonne, par exemple, de name à full_name.
  • Colonnes supprimées : Suppression de colonnes du schéma de table.
  • Élargissement de type : modification du type d'une colonne vers un type plus large. Par exemple, un champ INT devenant DOUBLE.
  • Autres changements de type : Changement du type d'une colonne. Par exemple, un champ INT devenant STRING.

Le support de l'évolution des schémas est essentiel pour construire des pipelines résilients et de longue durée qui peuvent s'adapter aux données changeantes sans mises à jour manuelles fréquentes.

Composants

L'évolution des schémas Databricks implique quatre principales catégories de composants, chacune gérant les modifications de schéma indépendamment :

  1. Connecteurs : Composants qui ingèrent des données provenant de sources externes. Ceux-ci incluent les connecteurs Auto Loader, Kafka, Kinesis et Lakeflow.
  2. Analyseurs de format : fonctions qui décodent les formats bruts, y compris from_json, from_avro, from_xml et from_protobuf.
  3. Moteurs : moteurs de traitement qui exécutent des queries, y compris Structured Streaming.
  4. **Datasets** : Tables de streaming, vues matérialisées, tables Delta et vues qui persistent et servent des données.

évolution des schémas

Chaque composant de l'architecture de Data Engineering de l'évolution des schémas est indépendant. Vous êtes responsable de la configuration de l'évolution des schémas dans les composants individuels afin d'obtenir le comportement souhaité dans votre flux de traitement de données.

Par exemple, lors de l'utilisation d'Auto Loader pour ingérer des données dans une table Delta, il existe deux schémas persistants : l'un est géré par Auto Loader à son emplacement de schéma et l'autre est le schéma de la table Delta cible. Dans un état stable, les deux sont identiques. Lorsque Auto Loader fait évoluer son schéma, en fonction des données entrantes, la table Delta doit également faire évoluer son schéma, sinon la query échoue. Dans ce cas, vous pouvez (a) mettre à jour le schéma de la table Delta cible en activant l'évolution des schémas ou en utilisant une commande DDL directe, ou (b) effectuer une réécriture complète de la table Delta cible.

Prise en charge de l'évolution des schémas par connecteur

Les sections suivantes détaillent comment chaque composant Databricks gère les différents types de modifications de schéma.

Auto Loader

Auto Loader prend en charge les modifications de colonne et l'élargissement de type. Configurez l'évolution automatique des schémas avec cloudFiles.schemaEvolutionMode et rescuedDataColumn. Vous pouvez définir manuellement schemaHints ou un schema immuable. Lorsque le schéma évolue automatiquement, le Stream échoue initialement. Au redémarrage, le schéma évolué est utilisé. Voir Comment fonctionne l'évolution des schémas Auto Loader ?.

  • Nouvelles colonnes : pris en charge, selon le schemaEvolutionMode sélectionné. Échec nécessitant un redémarrage manuel pour ajouter de nouvelles colonnes au schéma.
  • Renommage de colonne : Pris en charge, selon le schemaEvolutionMode sélectionné. La colonne renommée est traitée comme une nouvelle colonne ajoutée, et l'ancienne colonne est remplie avec NULL pour les nouvelles lignes. Échec avec un redémarrage manuel nécessaire pour mettre à jour le schéma.
  • Colonnes supprimées : pris en charge. Traitées comme des suppressions logiques, où les nouvelles lignes de la colonne supprimée sont définies sur NULL.
  • **Élargissement de type** : pris en charge dans Databricks Runtime 16.4 et versions ultérieures avec schemaEvolutionMode défini addNewColumnsWithTypeWidening sur. Les modifications prises en charge du type de données sont élargies automatiquement. Les modifications de type non prises en charge sont capturées dans le rescuedDataColumn. Consultez L’élargissement automatique du type avec Auto Loader.
  • Autres modifications de type : Non pris en charge. Les modifications de type sont capturées dans le rescuedDataColumn si rescueDataColumn a été défini et schemaEvolutionMode défini sur rescue. Sinon, cela nécessite un changement de schéma manuel.

Connecteur Delta

Le connecteur Delta peut prendre en charge l'évolution des schémas. Si vous lisez à partir d'une table Delta avec le mappage de colonnes et l'activation du suivi de schéma, cela prend en charge l'évolution des schémas pour le renommage et la suppression de colonnes. Vous devez définir la configuration Spark correcte pour chacun de ces changements respectifs afin de faire évoluer le schéma sans arrêter le Stream. Autrement, le Stream fait évoluer son schéma suivi dès qu'une modification est détectée, puis il s'arrête. Vous devez ensuite redémarrer manuellement la query de streaming pour reprendre le traitement.

  • Nouvelles colonnes : prises en charge. Lorsque mergeSchema est activé, les nouvelles colonnes sont ajoutées automatiquement. Autrement, la query échoue et vous devez redémarrer le Stream pour ajouter les nouvelles colonnes au schéma, mais la table Delta ne nécessite pas de réécriture.
  • **Renommage de colonne** : Pris en charge. Vous pouvez faire évoluer le schéma dans une query de streaming avec la configuration Spark spark.databricks.delta.streaming.allowSourceColumnRename.
  • Colonnes supprimées : pris en charge. Vous pouvez faire évoluer le schéma dans une query de streaming avec la configuration Spark spark.databricks.delta.streaming.allowSourceColumnDrop.
  • Élargissement de type : pris en charge dans Databricks Runtime 16.4 LTS et versions ultérieures. Lorsque mergeSchema est activé et que l'élargissement de type est activé sur la table cible, les modifications de type sont gérées automatiquement. Vous pouvez activer l’élargissement de type avec la propriété de table delta.enableTypeWidening. Voir Élargissement de type.
  • Autres changements de type : non pris en charge.

Connecteurs SaaS et CDC

Les connecteurs SaaS et CDC font évoluer automatiquement le schéma lorsque les colonnes changent. Ceci est géré par un redémarrage automatique lorsqu'un changement est détecté. Les modifications de type nécessitent un full refresh.

  • Nouvelles colonnes : prises en charge. La query redémarre automatiquement pour résoudre l'incohérence de schéma.
  • **Renommage de colonne** : Pris en charge. La query redémarre automatiquement pour résoudre l'incompatibilité de schéma. La colonne renommée est traitée comme une nouvelle colonne ajoutée.
  • Colonnes supprimées : pris en charge. Les colonnes supprimées sont traitées comme des suppressions logiques, où les nouvelles lignes pour la colonne supprimée sont définies sur NULL.
  • Élargissement de type : non pris en charge. La mise à jour du schéma nécessite un full refresh.
  • Autres modifications de type : Non pris en charge. La mise à jour du schéma nécessite un full refresh.

Connecteurs Kinesis, Kafka, Pub/Sub et Pulsar

Aucune évolution des schémas native n'est prise en charge. Chacune des fonctions de connecteur renvoie un blob binaire. L'évolution des schémas est gérée par le parseur de format.

  • Nouvelles colonnes : Traitées par l'analyseur de format.
  • Renommage de colonne : Géré par l’analyseur de format.
  • Colonnes ignorées : Traitées par l'analyseur de format.
  • Élargissement de type : Géré par l'analyseur de format.
  • **Autres changements de type** : Géré par l'analyseur de format.

Prise en charge de l'évolution des schémas par l'analyseur de format

from_json analyseur

L'analyseur from_json ne prend pas en charge l'évolution des schémas. Vous devez mettre à jour le schéma manuellement. Lorsque vous utilisez from_json dans les LakeFlow pipelines, l'évolution automatique des schémas peut être activée avec schemaLocationKey et schemaEvolutionMode.

  • Nouvelles colonnes : lorsque l'évolution automatique des schémas est activée, elle se comporte comme Auto Loader.
  • Renommage des colonnes : Lorsque l'évolution automatique des schémas est activée, elle se comporte comme Auto Loader.
  • Colonnes supprimées : lorsque l’évolution automatique des schémas est activée, elle se comporte comme Auto Loader.
  • Élargissement de type : lorsque l'évolution automatique des schémas est activée, elle se comporte comme Auto Loader.
  • **Autres modifications de type** : lorsque l'évolution automatique des schémas est activée, elle se comporte comme Auto Loader.

Analyseursfrom_avro et from_protobuf

Les analyseurs from_avro et from_protobuf se comportent de la même manière. Le schéma peut être récupéré depuis le Confluent Schema Registry, ou l'utilisateur peut fournir un schéma et doit le mettre à jour manuellement. Il n'y a pas de concept d'évolution des schémas au sein de la fonction from_avro ou from_protobuf ; elle doit être gérée par le moteur d'exécution et le Schema Registry.

  • Nouvelles colonnes : Prise en charge avec Confluent Schema Registry. Sinon, l'utilisateur doit mettre à jour le schéma manuellement.
  • Renommage de colonne : Pris en charge avec Confluent Schema Registry. Sinon, l'utilisateur doit mettre à jour le schéma manuellement.
  • Colonnes supprimées : pris en charge avec Confluent Schema Registry. Sinon, l'utilisateur doit mettre à jour le schéma manuellement.
  • Élargissement de type : Pris en charge avec Confluent Schema Registry. Sinon, l'utilisateur doit mettre à jour le schéma manuellement.
  • Autres modifications de type : Pris en charge avec Confluent Schema Registry. Sinon, l'utilisateur doit mettre à jour le schéma manuellement.

Analyseursfrom_csv et from_xml

Les analyseurs from_csv et from_xml ne prennent pas en charge l'évolution des schémas.

  • Nouvelles colonnes : Non pris en charge
  • Renommage de colonne : Non pris en charge
  • Colonnes supprimées : Non pris en charge
  • Élargissement de type : Non pris en charge
  • Autres modifications de type : Non pris en charge

Prise en charge de l'évolution des schémas par le moteur

Structured Streaming

Le schéma d'une query de streaming est verrouillé pendant la phase de planification, et tous les micro-batchs réutilisent ce plan sans replanification. Si le schéma source change en cours d'exécution, la query échoue et l'utilisateur doit redémarrer la query de streaming afin que Spark puisse replanifier en fonction du nouveau schéma.

Le dataset sur lequel le Stream écrit doit également prendre en charge l'évolution des schémas.

  • Nouvelles colonnes : prises en charge. La query échoue et vous devez redémarrer le Stream pour résoudre l'incompatibilité de schéma.
  • **Renommage de colonne** : Pris en charge. La query échoue et vous devez redémarrer le Stream pour résoudre l'incompatibilité de schéma.
  • Colonnes supprimées : pris en charge. La query échoue et vous devez redémarrer le Stream pour résoudre l'incompatibilité de schéma.
  • **Élargissement de type** : Pris en charge. La query échoue et vous devez redémarrer le Stream pour résoudre l'incompatibilité de schéma.
  • Autres modifications de type : Pris en charge. La query échoue et vous devez redémarrer le Stream pour résoudre l'incompatibilité de schéma.

Évolution des schémas par dataset

Tables de streaming

Les tables de streaming prennent en charge le comportement de Merge avec évolution des schémas par default. La mise à jour du schéma ne nécessite pas de redémarrage manuel, mais les modifications arbitraires du schéma nécessitent un refresh complet.

  • Nouvelles colonnes : prises en charge. La query redémarre automatiquement pour résoudre l'incompatibilité de schéma.
  • **Renommage de colonne** : Pris en charge. La query redémarre pour résoudre l'incompatibilité de schéma. La colonne renommée est traitée comme une nouvelle colonne ajoutée.
  • Colonnes supprimées : pris en charge. Les colonnes supprimées sont traitées comme des suppressions logiques, où les nouvelles lignes pour la colonne supprimée sont définies sur NULL.
  • **Élargissement de type** : Pris en charge. L'élargissement des types doit être activé soit au niveau du pipeline, soit directement sur la table. Voir l'élargissement des types dans les LakeFlow Pipelines.
  • Autres modifications de type : Non pris en charge. La mise à jour du schéma nécessite un full refresh.

Vues matérialisées

Toute mise à jour du schéma ou de la query de définition Trigger un recalcul complet de la vue matérialisée.

  • Nouvelles colonnes : Recalcul complet déclenché.
  • Renommage de colonne : nouveau calcul complet déclenché.
  • colonnes supprimées : nouveau calcul complet déclenché.
  • Élargissement du type : recalcul complet Trigger.
  • **Autres changements de type** : recalcul complet Trigger.

Tables Delta

Les tables Delta prennent en charge une variété de configurations pour mettre à jour le schéma de table, y compris le renommage, la suppression et l'élargissement du type de colonnes sans réécrire les données de la table. Les configurations prises en charge incluent l'évolution des schémas de fusion, le mappage de colonnes, l'élargissement des types et l'écrasement du schéma.

  • Nouvelles colonnes : prises en charge. Évolution automatique lorsque l'évolution des schémas Merge est activée, sans nécessiter de réécriture de table Delta. Si l'évolution des schémas Merge n'est pas activée, les mises à jour échouent.
  • **Renommage de colonne** : Pris en charge. Peut renommer via des commandes ALTER TABLE DDL manuelles avec le mappage de colonnes activé. Ne nécessite pas de réécriture de table Delta.
  • Colonnes supprimées : pris en charge. Peut supprimer des colonnes via des commandes ALTER TABLE DDL manuelles avec le mappage de colonnes activé. Ne nécessite pas de réécriture de table Delta.
  • **Élargissement de type** : Pris en charge. Applique automatiquement le changement de type lorsque l'élargissement de type et l'évolution des schémas de Merge sont activés. Vous pouvez élargir les colonnes via des commandes ALTER TABLE DDL manuelles lorsque l'élargissement de type est activé. Si aucune des deux n'est configurée, les opérations échouent. Voir Élargir les types avec l'évolution automatique des schémas.
  • Autres modifications de type : Prises en charge, mais nécessitent une réécriture complète de la table Delta. Vous devez activer overwriteSchema, ce qui permet une réécriture complète de la table Delta. Autrement, les opérations échouent.

Vues

Si la vue a un column_list qui ne correspond pas au nouveau schéma, ou si elle a une query qui ne peut pas être analysée, la vue devient invalide. Si ce n'est pas le cas, vous pouvez activer l'évolution des schémas pour les modifications de type avec SCHEMA TYPE EVOLUTION et pour les modifications de type, ainsi que pour les colonnes nouvelles, renommées et supprimées avec SCHEMA EVOLUTION (qui est un sur-ensemble de l'évolution de type).

  • Nouvelles colonnes : prises en charge. Avec le mode SCHEMA EVOLUTION, la vue évolue automatiquement sans aucune intervention manuelle s'il n'y a pas de column_list explicite. Sinon, la vue peut devenir non valide et l'utilisateur ne peut pas l'interroger.
  • Renommage des colonnes : Pris en charge. Avec le mode SCHEMA EVOLUTION, la vue évolue automatiquement sans aucune intervention manuelle s'il n'y a pas de column_list explicite. Sinon, la vue peut devenir invalide.
  • Colonnes supprimées : pris en charge. Avec le mode SCHEMA EVOLUTION, la vue évolue automatiquement sans aucune intervention manuelle s'il n'y a pas de column_list explicite. Sinon, la vue peut devenir invalide.
  • **Élargissement de type** : Pris en charge. Avec le mode SCHEMA TYPE EVOLUTION, la vue évolue automatiquement pour toute modification de type. Avec le mode SCHEMA EVOLUTION, la vue évolue automatiquement sans aucune intervention manuelle s’il n’y a pas de column_list explicite. Sinon, la vue peut devenir invalide.
  • Autres modifications de type : Pris en charge. Avec le mode SCHEMA TYPE EVOLUTION, la vue évolue automatiquement pour toute modification de type. Avec le mode SCHEMA EVOLUTION, la vue évolue automatiquement sans aucune intervention manuelle s’il n’y a pas de column_list explicite. Sinon, la vue peut devenir invalide.

Exemple

L'exemple suivant montre comment ingérer un sujet Kafka avec des charges utiles encodées en Avro enregistrées dans Confluent Schema Registry, et les écrire dans une table Delta gérée avec l'évolution des schémas activée.

Points clés illustrés :

  • Intégrer avec le connecteur Kafka.
  • Décodez les enregistrements Avro en utilisant from_avro avec un registre de schémas Kafka.
  • Gérer l'évolution des schémas en définissant avroSchemaEvolutionMode.
  • Écrivez dans une table Delta avec mergeSchema activé pour autoriser les modifications additives.

Le code suppose que vous disposez d'un sujet Kafka utilisant un registre de schémas Confluent, et qui émet des données encodées en Avro.

Python
# ----- CONFIG: fill these in -----
# Catalog and schema:
CATALOG = "<catalog_name>"
SCHEMA = "<schema_name>"
# Schema Registry:
# (This is where the producer evolves the schema)
SCHEMA_REG = "<schema registry endpoint>"
SR_USER = "<api key>"
SR_PASS = "<api secret>"
# Confluent Cloud: SASL_SSL broker:
BOOTSTRAP = "<server:ip>"
# Kafka topic:
TOPIC = "<topic>"
# ----- end: config -----

BRONZE_TABLE = f"{CATALOG}.{SCHEMA}.bronze_users"
CHECKPOINT = f"/Volumes/{CATALOG}/{SCHEMA}/checkpoints/bronze_users"

# Kafka auth (example for Confluent Cloud SASL/PLAIN over SSL)
KAFKA_OPTS = {
"kafka.security.protocol": "SASL_SSL",
"kafka.sasl.mechanism": "PLAIN",
"kafka.sasl.jaas.config": f"kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username='{SR_USER}' password='{SR_PASS}';"
}

# ----- Evolution knobs -----
# spark.conf.set("spark.databricks.delta.schema.autoMerge.enabled", value = True)

from pyspark.sql.functions import col
from pyspark.sql.avro.functions import from_avro

# Build reader
reader = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", BOOTSTRAP)
.option("subscribe", TOPIC)
.option("startingOffsets", "earliest")
)

# Attach Kafka auth options
for k, v in KAFKA_OPTS.items():
reader = reader.option(k, v)

# --- No native schema evolution supported. Returns a binary blob. ---
raw_df = reader.load()

# Decode Avro with Schema Registry
# --- The format parser handles updating the schema using the schema registry ---
decoded = from_avro(
data=col("value"),
jsonFormatSchema=None, # using SR
subject=f"{TOPIC}-value",
schemaRegistryAddress=SCHEMA_REG,
options={
&quot;confluent.schema.registry.basic.auth.credentials.source&quot;: &quot;USER_INFO&quot;,
&quot;confluent.schema.registry.basic.auth.user.info&quot;: f&quot;{SR_USER}:{SR_PASS}&quot;,
# Behavior on schema changes:
&quot;avroSchemaEvolutionMode&quot;: &quot;restart&quot;, # fail-fast so you can restart and adopt new fields
&quot;mode&quot;: &quot;FAILFAST&quot;
}
).alias("payload")

bronze_df = raw_df.select(decoded, "timestamp").select("payload.*", "timestamp")

# Write to a managed Delta table as a STREAM
# --- Need to enable schema evolution separately for streaming to a Delta separately with mergeSchema --
(bronze_df.writeStream
.format("delta")
.option("checkpointLocation", CHECKPOINT)
.option("ignoreChanges", "true")
.outputMode("append")
.option("mergeSchema", "true") # only supports adding new columns. Renaming, dropping, and type changes need to be handled separately.
.trigger(availableNow=True) # Use availableNow trigger for Databricks SQL/Unity Catalog
.toTable(BRONZE_TABLE)
)