Aller au contenu principal

Déplacer les tables entre les pipelines

Déplacez les tables de streaming et les vues matérialisées entre les pipelines afin que le pipeline de destination mette à jour la table au lieu de l'originale. Ceci est utile dans de nombreux scénarios, notamment :

  • Divisez un grand pipeline en de plus petits.
  • Merge plusieurs pipelines en un seul plus grand.
  • Modifiez la fréquence de refresh de certaines tables dans un pipeline.
  • Déplacez les tables d'un pipeline qui utilise le mode de publication hérité vers le mode de publication default. Pour plus de détails sur le mode de publication hérité, consultez le mode de publication hérité pour les pipelines. Pour savoir comment vous pouvez migrer le mode de publication pour un pipeline entier en une seule fois, consultez Activer le mode de publication par default dans un pipeline.
  • Déplacez les tables entre les pipelines dans différents Workspaces.

Exigences

Voici les exigences pour le déplacement d'une table entre les pipelines.

  • Vous devez utiliser Databricks Runtime 16.3 ou version ultérieure lors de l'exécution de la commande ALTER ..., et Databricks Runtime 17.2 pour le déplacement de table inter-Workspace.

  • Les pipelines source et de destination doivent se trouver dans des workspaces qui partagent un metastore. Pour vérifier le metastore, consultez la fonction current_metastore.

  • Le compte utilisateur ou le Service Principal exécutant l'opération doit être l'utilisateur « Exécuter en tant que » des pipelines source et de destination.

  • Vous devez être le propriétaire des pipelines source et de destination.

  • Le pipeline de destination doit utiliser le mode de publication default. Cela vous permet de publier des tables sur plusieurs catalogues et schémas.

    Alternativement, les deux pipelines doivent utiliser le mode de publication hérité et les deux doivent avoir le même catalogue et la même valeur cible dans les paramètres. Pour des informations sur le mode de publication hérité, consultez schéma LIVE (hérité).

    Les pipelines en mode de publication hérité sont indiqués dans le champ **Récapitulatif** de l'interface utilisateur des paramètres du pipeline.

remarque

Cette fonctionnalité ne prend pas en charge le déplacement d'un pipeline utilisant le mode de publication default vers un pipeline utilisant le mode de publication hérité.

Déplacer une table entre les pipelines

Les instructions suivantes décrivent comment déplacer une table de streaming ou une vue matérialisée d’un pipeline à un autre.

  1. Arrêtez le pipeline source s'il est en cours d'exécution. Attendez qu'il s'arrête complètement.

  2. Supprimez la définition de la table du code source du pipeline et stockez-la quelque part pour référence ultérieure.

    Incluez toutes les requêtes ou le code de support nécessaires au bon fonctionnement du pipeline.

  3. Depuis un Notebook ou un éditeur SQL, exécutez la commande SQL suivante pour réaffecter la table du pipeline source au pipeline de destination :

    SQL
    ALTER [MATERIALIZED VIEW | STREAMING TABLE | TABLE] <table-name>
    SET TBLPROPERTIES("pipelines.pipelineId"="<destination-pipeline-id>");

    Notez que la commande SQL doit être exécutée depuis le workspace du pipeline source.

    La commande utilise ALTER MATERIALIZED VIEW et ALTER STREAMING TABLE pour les vues matérialisées gérées et les tables de streaming Unity Catalog, respectivement. Pour effectuer la même action sur une table Hive metastore, utilisez ALTER TABLE.

    Par exemple, si vous souhaitez déplacer une table de streaming nommée sales vers un pipeline avec l'ID abcd1234-ef56-ab78-cd90-1234efab5678, vous exécuteriez la commande suivante :

    SQL
    ALTER STREAMING TABLE sales
    SET TBLPROPERTIES("pipelines.pipelineId"="abcd1234-ef56-ab78-cd90-1234efab5678");
remarque

Le pipelineId doit être un identifiant de pipeline valide. La valeur null n'est pas autorisée.

  1. Ajoutez la définition de la table au code du pipeline de destination.
remarque

Si le catalogue ou le schéma cible diffèrent entre la source et la destination, copier la query exactement pourrait ne pas fonctionner. Les tables partiellement qualifiées dans la définition peuvent se résoudre différemment. Vous devrez peut-être mettre à jour la définition lors du déplacement pour qualifier entièrement les noms des tables.

remarque

Supprimez ou commentez tout flux d'ajout à usage unique (en Python, les requêtes avec append_flow(once=True); en SQL, les requêtes avec INSERT INTO ONCE) du code du pipeline de destination. Pour plus de détails, consultez les Limitations.

Le déplacement est terminé. Vous pouvez maintenant exécuter les pipelines source et de destination. Le pipeline de destination met à jour la table.

Dépannage

Le tableau suivant décrit les erreurs qui pourraient survenir lors du déplacement d'une table entre les pipelines.

Erreur

Description

DESTINATION_PIPELINE_NOT_IN_DIRECT_PUBLISHING_MODE

Le pipeline source est en mode de publication default, et la destination utilise le mode de schéma LIVE (hérité). Ceci n'est pas pris en charge. Si la source utilise le mode de publication default, la destination doit également le faire.

PIPELINE_TYPE_NOT_WORKSPACE_PIPELINE_TYPE

Seul le déplacement de tables entre les pipelines est pris en charge. Le déplacement des tables de streaming autonomes et des vues matérialisées n’est pas pris en charge.

DESTINATION_PIPELINE_NOT_FOUND

Le pipelines.pipelineId doit être un pipeline valide. Le pipelineId ne peut pas être nul.

Échec de la mise à jour de la table dans la destination après le déplacement.

Pour rapidement atténuer (le problème) dans ce cas, redéplacez la table vers le pipeline source en suivant les mêmes instructions.

PIPELINE_PERMISSION_DENIED_NOT_OWNER

Les pipelines source et de destination doivent appartenir à l'utilisateur effectuant l'opération de déplacement.

TABLE_ALREADY_EXISTS

La table mentionnée dans le message d'erreur existe déjà. Cela peut arriver si une table de support pour le pipeline existe déjà. Dans ce cas, DROP la table mentionnée dans l'erreur.

Erreur

Description

DESTINATION_PIPELINE_NOT_IN_DIRECT_PUBLISHING_MODE

Le pipeline source est en mode de publication default, et la destination utilise le mode de schéma LIVE (hérité). Ceci n'est pas pris en charge. Si la source utilise le mode de publication default, la destination doit également le faire.

PIPELINE_TYPE_NOT_WORKSPACE_PIPELINE_TYPE

Seul le déplacement de tables entre les pipelines est pris en charge. Le déplacement des tables de streaming autonomes et des vues matérialisées n’est pas pris en charge.

DESTINATION_PIPELINE_NOT_FOUND

Le pipelines.pipelineId doit être un pipeline valide. Le pipelineId ne peut pas être nul.

Échec de la mise à jour de la table dans la destination après le déplacement.

Pour rapidement atténuer (le problème) dans ce cas, redéplacez la table vers le pipeline source en suivant les mêmes instructions.

PIPELINE_PERMISSION_DENIED_NOT_OWNER

Les pipelines source et de destination doivent appartenir à l'utilisateur effectuant l'opération de déplacement.

TABLE_ALREADY_EXISTS

La table mentionnée dans le message d'erreur existe déjà. Cela peut arriver si une table de support pour le pipeline existe déjà. Dans ce cas, DROP la table mentionnée dans l'erreur.

Exemple avec plusieurs tables dans un pipeline

Les pipelines peuvent contenir plus d'une table. Vous pouvez toujours déplacer une seule table à la fois entre les pipelines. Dans ce scénario, il y a trois tables (table_a, table_b, table_c) qui lisent les unes des autres séquentiellement dans le pipeline source. Nous voulons déplacer une table, table_b, vers un autre pipeline.

Code source initial du pipeline :

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table
def table_a():
return spark.read.table("source_table")

# Table to be moved to new pipeline:
@dp.table
def table_b():
return (
spark.read.table("table_a")
.select(col("column1"), col("column2"))
)

@dp.table
def table_c():
return (
spark.read.table("table_b")
.groupBy(col("column1"))
.agg(sum("column2").alias("sum_column2"))
)

Nous déplaçons table_b vers un autre pipeline en copiant et en supprimant la définition de table de la source et en mettant à jour le pipelineId de table_b.

Tout d'abord, suspendez les planifications et attendez que les mises à jour soient terminées sur les pipelines source et cible. Modifiez ensuite le pipeline source pour supprimer le code de la table en cours de déplacement. Le code d'exemple de pipeline source mis à jour devient :

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table
def table_a():
return spark.read.table("source_table")

# Removed, to be in new pipeline:
# @dp.table
# def table_b():
# return (
# spark.read.table("table_a")
# .select(col("column1"), col("column2"))
# )

@dp.table
def table_c():
return (
spark.read.table("table_b")
.groupBy(col("column1"))
.agg(sum("column2").alias("sum_column2"))
)

Accédez à l'éditeur SQL pour exécuter la commande ALTER pipelineId.

SQL
ALTER MATERIALIZED VIEW table_b
SET TBLPROPERTIES("pipelines.pipelineId"="<new-pipeline-id>");

Ensuite, accédez au pipeline de destination et ajoutez la définition de table_b. Si le catalogue et le schéma default sont les mêmes dans les paramètres du pipeline, aucune modification de code n’est requise.

Le code du pipeline cible :

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table(name="table_b")
def table_b():
return (
spark.read.table("table_a")
.select(col("column1"), col("column2"))
)

Si le catalogue et le schéma par default diffèrent dans les paramètres du pipeline, vous devez ajouter le nom entièrement qualifié en utilisant le catalogue et le schéma du pipeline.

Par exemple, le code du pipeline cible pourrait être :

Python
from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.table(name="source_catalog.source_schema.table_b")
def table_b():
return (
spark.read.table("source_catalog.source_schema.table_a")
.select(col("column1"), col("column2"))
)

Exécutez (ou réactivez les planifications) pour les pipelines source et cible.

Les pipelines sont maintenant disjoints. La query pour table_c lit à partir de table_b (maintenant dans le pipeline cible) et table_b lit à partir de table_a (dans le pipeline source). Lorsque vous effectuez une exécution déclenchée sur le pipeline source, table_b n'est pas mis à jour car il n'est plus géré par le pipeline source. Le pipeline source traite table_b comme une table externe au pipeline. Ceci est comparable à la définition d'une vue matérialisée lisant à partir d'une table Delta dans Unity Catalog qui n'est pas gérée par le pipeline.

Limitations

Voici les limitations pour le déplacement de tables entre les pipelines.

  • Les vues matérialisées autonomes et les tables de streaming ne sont pas prises en charge.
  • Les flux d'ajout unique – les flux Python append_flow(once=True) et les flux SQL INSERT INTO ONCE – ne sont pas pris en charge. Leurs états d'exécution ne sont pas conservés et ils peuvent s'exécuter à nouveau dans le pipeline de destination. Supprimez ou commentez les flux d'ajout unique du pipeline de destination pour éviter d'exécuter à nouveau ces flux.
  • Les tables ou vues privées ne sont pas prises en charge.
  • Les pipelines source et de destination doivent être des pipelines. Les pipelines nuls ne sont pas pris en charge.
  • Les pipelines source et de destination doivent être soit dans le même Workspace, soit dans des Workspaces différents qui partagent le même métastore.
  • Les pipelines source et de destination doivent appartenir à l'utilisateur qui exécute l'opération de déplacement.
  • Si le pipeline source utilise le mode de publication default, le pipeline de destination doit également utiliser le mode de publication default. Vous ne pouvez pas déplacer une table d’un pipeline utilisant le mode de publication default vers un pipeline qui utilise le schéma LIVE (hérité). Consultez le schéma LIVE (hérité).
  • Si les pipelines source et de destination utilisent tous les deux le schéma LIVE (hérité), alors ils doivent avoir les mêmes valeurs catalog et target dans les paramètres.