Aller au contenu principal

Créer un pipeline Unity Catalog en clonant un pipeline Hive metastore

La requête clone a pipeline dans l'API REST Databricks copie un pipeline existant qui publie dans le Hive metastore vers un nouveau pipeline qui publie dans Unity Catalog. Lorsque vous appelez la requête clone a pipeline, celle-ci :

  • Copie le code source et la configuration du pipeline existant vers un nouveau, en appliquant toutes les surcharges de configuration que vous avez spécifiées.
  • Met à jour les définitions et les références des vues matérialisées et des tables de streaming avec les modifications requises pour que ces objets soient gérés par Unity Catalog.
  • Démarre une mise à jour de pipeline pour migrer les données et métadonnées existantes, telles que les points de contrôle, pour toute table de streaming dans le pipeline. Cela permet à ces tables de streaming de reprendre le traitement au même point que le pipeline d'origine.

Une fois l'opération de clonage terminée, les pipelines d'origine et nouveaux peuvent s'exécuter indépendamment.

Cette page contient des exemples d'appel direct de la requête API et via un script Python à partir d'un Notebook Databricks.

Avant de commencer

Les éléments suivants sont requis avant de cloner un pipeline :

  • Pour cloner un pipeline Hive metastore, les tables et vues définies dans le pipeline doivent publier des tables vers un schéma cible. Pour savoir comment ajouter un schéma cible à une pipeline, consultez Configurer un pipeline pour publier vers Hive metastore.

  • Les références aux tables ou vues gérées par le Hive metastore dans le pipeline à cloner doivent être entièrement qualifiées avec le catalogue (hive_metastore), le schéma et le nom de la table. Par exemple, dans le code suivant qui crée un customers dataset, l'argument du nom de table doit être mis à jour en hive_metastore.sales.customers:

    Python
    @dp.table
    def customers():
    return spark.read.table("sales.customers").where(...)
  • Ne modifiez pas le code source du pipeline Hive metastore pendant qu’une opération de clonage est en cours, y compris les Notebooks configurés dans le pipeline et tous les modules stockés dans les dossiers Git ou les fichiers du Workspace.

  • Le pipeline source du Hive metastore ne doit pas être en cours d'exécution lorsque vous start l'opération de clonage. Si une mise à jour est en cours d'exécution, arrêtez-la ou attendez qu'elle se termine.

Voici d'autres considérations importantes avant de cloner un pipeline :

  • Si les tables du pipeline Hive metastore spécifient un emplacement de stockage à l'aide de l'argument path dans Python ou LOCATION dans SQL, transmettez la configuration "pipelines.migration.ignoreExplicitPath": "true" à la demande de clonage. La définition de cette configuration est incluse dans les instructions ci-dessous.
  • Si le pipeline du métastore Hive inclut une source Auto Loader qui spécifie une valeur pour l'option cloudFiles.schemaLocation, et que le pipeline du métastore Hive reste opérationnel après la création du clone Unity Catalog, vous devez définir l'option mergeSchema sur true à la fois dans le pipeline du métastore Hive et dans le pipeline Unity Catalog cloné. L'ajout de cette option au pipeline du Hive metastore avant le clonage copie l'option vers le nouveau pipeline.

Cloner un pipeline avec l'API REST Databricks

L'exemple suivant utilise la commande curl pour appeler la requête clone a pipeline dans l'API REST Databricks :

Bash
curl -X POST \
--header "Authorization: Bearer <personal-access-token>" \
<databricks-instance>/api/2.0/pipelines/<pipeline-id>/clone \
--data @clone-pipeline.json

Remplacer :

clone-pipeline.json :

JSON
{
"catalog": "<target-catalog-name>",
"target": "<target-schema-name>",
"name": "<new-pipeline-name>",
"clone_mode": "MIGRATE_TO_UC",
"configuration": {
"pipelines.migration.ignoreExplicitPath": "true"
}
}

Remplacer :

  • <target-catalog-name> avec le nom d'un catalogue dans Unity Catalog vers lequel le nouveau pipeline doit publier. Cela doit être un catalogue existant.
  • <target-schema-name> avec le nom d'un schéma dans Unity Catalog vers lequel le nouveau pipeline doit publier s'il est différent du nom de schéma actuel. Ce paramètre est facultatif et, s’il n’est pas spécifié, le nom de schéma existant est utilisé.
  • <new-pipeline-name> avec un nom facultatif pour le nouveau pipeline. Si non spécifié, le nouveau pipeline est nommé en utilisant le nom du pipeline source suivi de [UC].

clone_mode spécifie le mode à utiliser pour l'opération de clonage. MIGRATE_TO_UC est la seule option prise en charge.

Utilisez le champ configuration pour spécifier les configurations sur le nouveau pipeline. Les valeurs définies ici remplacent les configurations dans le pipeline d'origine.

La réponse de la requête d'API REST clone est l'ID de pipeline du nouveau pipeline Unity Catalog.

Cloner un pipeline à partir d'un notebook Databricks

L'exemple suivant appelle la requête create a pipeline à partir d'un script Python. Vous pouvez utiliser un Notebook Databricks pour exécuter ce script :

  1. Créez un nouveau Notebook pour le script. Voir Créer un Notebook.

  2. Copiez le script Python suivant dans la première cellule du Notebook.

  3. Mettez à jour les valeurs d'espace réservé dans le script en remplaçant :

    • <databricks-instance> avec le nom d'instance du Databricks Workspace, par exemple dbc-a1b2345c-d6e7.cloud.databricks.com
    • <pipeline-id> avec l'identifiant unique du pipeline Hive metastore à cloner. Vous trouverez l'ID du pipeline dans l'interface utilisateur des pipelines.
    • <target-catalog-name> avec le nom d'un catalogue dans Unity Catalog vers lequel le nouveau pipeline doit publier. Cela doit être un catalogue existant.
    • <target-schema-name> avec le nom d'un schéma dans Unity Catalog vers lequel le nouveau pipeline doit publier s'il est différent du nom de schéma actuel. Ce paramètre est facultatif et, s’il n’est pas spécifié, le nom de schéma existant est utilisé.
    • <new-pipeline-name> avec un nom facultatif pour le nouveau pipeline. Si non spécifié, le nouveau pipeline est nommé en utilisant le nom du pipeline source suivi de [UC].
  4. Exécutez le script. Consultez Exécuter des Notebooks Databricks.

Python
import requests

# Your Databricks workspace URL, with no trailing spaces
WORKSPACE = "<databricks-instance>"

# The pipeline ID of the Hive metastore pipeline to clone
SOURCE_PIPELINE_ID = "<pipeline-id>"
# The target catalog name in Unity Catalog
TARGET_CATALOG = "<target-catalog-name>"
# (Optional) The name of a target schema in Unity Catalog. If empty, the same schema name as the Hive metastore pipeline is used
TARGET_SCHEMA = "<target-schema-name>"
# (Optional) The name of the new pipeline. If empty, the following is used for the new pipeline name: f"{originalPipelineName} [UC]"
CLONED_PIPELINE_NAME = "<new-pipeline-name>"

# This is the only supported clone mode
CLONE_MODE = "MIGRATE_TO_UC"

# Specify override configurations
OVERRIDE_CONFIGS = {"pipelines.migration.ignoreExplicitPath": "true"}

def get_token():
ctx = dbutils.notebook.entry_point.getDbutils().notebook().getContext()
return getattr(ctx, "apiToken")().get()

def check_source_pipeline_exists():
data = requests.get(
f"{WORKSPACE}/api/2.0/pipelines/{SOURCE_PIPELINE_ID}",
headers={&quot;Authorization&quot;: f&quot;Bearer {get_token()}&quot;},
)

assert data.json()["pipeline_id"] == SOURCE_PIPELINE_ID, "The provided source pipeline does not exist!"

def request_pipeline_clone():
payload = {
"catalog": TARGET_CATALOG,
"clone_mode": CLONE_MODE,
}
if TARGET_SCHEMA != "":
payload["target"] = TARGET_SCHEMA
if CLONED_PIPELINE_NAME != "":
payload["name"] = CLONED_PIPELINE_NAME
if OVERRIDE_CONFIGS:
payload["configuration"] = OVERRIDE_CONFIGS

data = requests.post(
f"{WORKSPACE}/api/2.0/pipelines/{SOURCE_PIPELINE_ID}/clone",
headers={&quot;Authorization&quot;: f&quot;Bearer {get_token()}&quot;},
json=payload,
)
response = data.json()
return response

check_source_pipeline_exists()
request_pipeline_clone()

Limitations

Voici les limitations de la requête d’API clone a pipeline :

  • Le clonage d'un pipeline de Hive metastore vers Unity Catalog n'est pas pris en charge à l'aide des Bundles d'automatisation déclaratifs.

  • Seul le clonage d'un pipeline configuré pour utiliser le Hive metastore vers un pipeline Unity Catalog est pris en charge.

  • Vous pouvez créer un clone uniquement dans le même Workspace Databricks que le pipeline à partir duquel vous effectuez le clonage.

  • Le pipeline que vous clonez ne peut inclure que les sources de streaming suivantes :

    • Sources Delta
    • Auto Loader, y compris toutes les sources de données prises en charge par Auto Loader. Voir Charger des fichiers depuis le stockage d'objets cloud.
    • Apache Kafka avec Structured Streaming. Cependant, la source Kafka ne peut pas être configurée pour utiliser l'option kafka.group.id. Consultez Connectez-vous à Apache Kafka.
    • Amazon Kinesis avec Structured Streaming. Cependant, la source Kinesis ne peut pas être configurée pour définir consumerMode sur efo.
  • Si le pipeline du Hive metastore que vous clonez utilise le mode de notification de fichiers Auto Loader, Databricks recommande de ne pas exécuter le pipeline du Hive metastore après le clonage. C'est parce que l'exécution du pipeline du Hive metastore entraîne la suppression de certains événements de notification de fichiers du clone Unity Catalog. Si le pipeline du Hive metastore source s'exécute une fois l'opération de clonage terminée, vous pouvez recharger les fichiers manquants à l'aide d'Auto Loader avec l'option cloudFiles.backfillInterval. Pour en savoir plus sur le mode de notification de fichiers d’Auto Loader, consultez Configurer les flux Auto Loader en mode de notification de fichiers. Pour en savoir plus sur le remplissage de fichiers avec Auto Loader, consultez Trigger des remplissages réguliers à l’aide de cloudFiles.backfillInterval. et Courant.

  • Les tâches de maintenance du pipeline sont automatiquement mises en pause pour les deux pipelines pendant que le clonage est en cours.

  • Ce qui suit s'applique aux queries time travel sur les tables du pipeline Unity Catalog cloné :

    • Si une version de table a été initialement écrite dans un objet géré par Hive metastore, les query time travel utilisant une clause timestamp_expression sont indéfinies lors de l'interrogation de l'objet Unity Catalog cloné.
    • Cependant, si la version de la table a été écrite dans l'objet Unity Catalog cloné, les queries time travel utilisant une clause timestamp_expression fonctionnent correctement.
    • Les requêtes Time travel utilisant une clause version fonctionnent correctement lors de l'interrogation d'un objet Unity Catalog cloné, même lorsque la version a été initialement écrite dans l'objet géré du Hive metastore.
  • Pour les autres limitations lors de l’utilisation de LakeFlow Pipelines avec Unity Catalog, consultez les limitations des pipelines Unity Catalog.

  • Pour les limitations d'Unity Catalog, consultez Unity Catalog overview.