Mettez à jour les Job lorsque vous mettez à niveau les Workspace hérités vers Unity Catalog
Lorsque vous mettez à niveau les Workspaces hérités vers Unity Catalog, vous devrez peut-être mettre à jour les jobs existants pour référencer les tables et chemins de fichiers mis à niveau. Le tableau principal de cette page répertorie les scénarios et les suggestions typiques pour la mise à jour de vos Jobs. Les scénarios qui nécessitent des exemples de code renvoient à la section Scénarios détaillés.
Pour une démonstration de la mise à jour des Jobs vers Unity Catalog, consultez Mise à niveau d'un Job vers Unity Catalog.
Vue d'ensemble des scénarios
Scénario | Solutions |
|---|---|
Le Job utilise des bibliothèques personnalisées via un script d'initialisation ou des bibliothèques définies par le cluster. Une bibliothèque personnalisée serait définie comme un | Modifiez la bibliothèque personnalisée pour vous assurer que :
|
Un Job lit ou écrit dans une table du Hive metastore. |
|
Le Job lit ou écrit dans des chemins qui sont des sous-dossiers de tables (non pris en charge dans Unity Catalog). |
|
Le Job lit ou écrit dans des chemins de montage qui sont des tables Unity Catalog. |
|
Le Job lit ou écrit des fichiers (pas des tables) en utilisant des chemins de montage. | Modifiez le code pour écrire à la place dans un emplacement de volume. |
Job est un job de streaming qui utilise | Non pris en charge actuellement. Envisagez de réécrire si possible, ou n'essayez pas de refactoriser ce Job tant qu'aucun support n'est fourni. |
Le Job est un Job de streaming qui utilise le mode de traitement continu. | Le mode de traitement continu est expérimental dans Spark et n'est pas pris en charge dans Unity Catalog. Refactorisez le Job pour utiliser le streaming structuré. Si cela n'est pas possible, envisagez de maintenir le job en cours d'exécution par rapport à Hive Metastore. |
Le Job est un job de streaming qui utilise des répertoires de points de contrôle. |
|
Le Job a une définition de cluster inférieure à Databricks Runtime 11,3. |
|
Le Job contient des Notebooks qui interagissent avec le stockage ou les tables. | Le Service Principal avec lequel le Job a été exécuté doit disposer d'un accès en lecture et en écriture aux Ressources requises dans Unity Catalog, telles que les volumes, les tables, les emplacements externes, etc. |
Job est un pipeline LakeFlow. |
|
Job utilise des services cloud non liés au stockage (tels qu'AWS Kinesis) avec un profil d'instance pour l'authentification. | Modifiez le code pour utiliser les identifiants de service Unity Catalog, qui régissent les identifiants capables d'interagir avec les services cloud non liés au stockage en générant des identifiants temporaires utilisables par les SDK. |
Job utilise Scala. |
|
Le Job contient des Notebooks qui utilisent des UDF Scala. |
|
Le Job a des tâches qui utilisent MLR. | Exécutez sur un compute dédié. |
Le Job possède une configuration de clusters qui repose sur des scripts d’initialisation globaux. |
|
Le Job utilise des fichiers JAR/Maven, des extensions Spark ou des sources de données personnalisées (à partir de Spark). |
|
Le Job contient des Notebooks avec des UDF PySpark. | Utilisez Databricks Runtime 13,2 ou une version ultérieure. |
Le Job contient des Notebooks avec du code Python qui effectue des appels réseau. | Utilisez Databricks Runtime 12.2 ou une version ultérieure. |
Le Job contient des Notebooks avec des UDF Pandas (scalaires). | Utilisez Databricks Runtime 13,2 ou une version ultérieure. |
Job utilisera les volumes Unity Catalog. | Utilisez Databricks Runtime 13.3 ou une version ultérieure. |
Le Job utilise | Voir les notebooks Job utilisent spark.catalog.X sur un cluster partagé. |
Le Job utilise | Utilisez Nécessite Databricks Runtime 13.3 LTS+ |
Le Job utilise |
|
Le Job utilise | Consultez Les Notebook de Job utilisent sc.parallelize et spark.read.json() sur un cluster partagé. |
Le Job utilise | Consultez Job Notebooks créent des DataFrames vides à l'aide de sc.emptyRDD() sur un cluster partagé. |
Le Job utilise le RDD | Les clusters partagés Unity Catalog utilisent Spark Connect pour la communication entre les programmes Python et Scala et le Spark Server, rendant les RDD inaccessibles. Un cas d'utilisation typique pour les RDD est d'exécuter une logique d'initialisation coûteuse une seule fois, puis d'effectuer des opérations moins coûteuses par ligne, telles que l'appel d'un service externe ou l'initialisation de la logique de chiffrement. Réécrivez les opérations RDD à l'aide de l'API DataFrame et des UDF Arrow natifs de PySpark. |
Le Job utilise SparkContext ( |
La JVM Spark n’est pas directement accessible depuis les REPL Python ou Scala — uniquement via les commandes Spark. Les Les |
Le Job utilise |
|
Le Job utilise |
Utilisez Nécessite Databricks Runtime 14,1 ou une version ultérieure. |
Le Job utilise | Sur les clusters partagés, le Spark Context n'est pas accessible pour définir les niveaux de logs directement. Dans Databricks Runtime 14+, le Spark Context n'est plus disponible. Définissez |
Le Job utilise des expressions ou des queries profondément imbriquées sur un cluster partagé. | Les DataFrames profondément imbriqués et les expressions créées de manière récursive à l'aide de l'API PySpark DataFrame peuvent produire :
Identifiez les chemins de code profondément imbriqués et réécrivez-les à l'aide d'expressions linéaires, de sous-requêtes ou de vues temporaires. Par exemple, au lieu d'appeler de manière récursive |
Le Job utilise | Voir les Notebooks Job utilisent input_file_name() sur un cluster partagé. |
Le Job effectue des opérations de données sur DBFS sur un cluster partagé. | Consultez les Notebooks de Job effectuant des Opérations de données sur DBFS sur un cluster partagé. |
Scénarios détaillés
Les scénarios suivants nécessitent des exemples de code.
Les Notebooks Job utilisent spark.catalog.X sur un cluster partagé
Utiliser Databricks Runtime 14,2 ou une version ultérieure.
Si une mise à niveau de Databricks Runtime n'est pas possible, utilisez les solutions de contournement suivantes.
Au lieu de tableExists, utilisez :
# SQL workaround
def tableExistsSql(tablename):
try:
spark.sql(f"DESCRIBE TABLE {tablename};")
except Exception as e:
return False
return True
tableExistsSql("jakob.jakob.my_table")
Au lieu de listTables, utilisez SHOW TABLES (qui prend également en charge la restriction par base de données ou la correspondance de modèles) :
spark.sql("SHOW TABLES")
Pour setDefaultCatalog, exécutez :
spark.sql("USE CATALOG <catalog_name>")
Les Notebooks de Job utilisent sc.parallelize et spark.read.json() sur un cluster partagé
Utilisez json.loads à la place.
Avant :
json_content1 = "{'json_col1': 'hello', 'json_col2': 32}"
json_content2 = "{'json_col1': 'hello', 'json_col2': 'world'}"
json_list = []
json_list.append(json_content1)
json_list.append(json_content2)
df = spark.read.json(sc.parallelize(json_list))
display(df)
Après :
from pyspark.sql import Row
import json
# Sample JSON data as a list of dictionaries (similar to JSON objects)
json_data_str = response.text
json_data = [json.loads(json_data_str)]
# Convert dictionaries to Row objects
rows = [Row(**json_dict) for json_dict in json_data]
# Create DataFrame from list of Row objects
df = spark.createDataFrame(rows)
df.display()
Les Notebook de Job créent des DataFrames vides à l'aide de sc.emptyRDD() sur un cluster partagé
Avant :
val schema = StructType( StructField("k", StringType, true) :: StructField("v", IntegerType, false) :: Nil)
spark.createDataFrame(sc.emptyRDD[Row], schema)
Après :
import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType}
val schema = StructType( StructField("k", StringType, true) :: StructField("v", IntegerType, false) :: Nil)
spark.createDataFrame(new java.util.ArrayList[Row](), schema)
from pyspark.sql.types import StructType, StructField, StringType
schema = StructType([StructField("k", StringType(), True)])
spark.createDataFrame([], schema)
Les Notebooks Job utilisent input_file_name() sur un cluster partagé
input_file_name() n'est pas pris en charge dans Unity Catalog pour les clusters partagés.
Pour obtenir le nom du fichier :
.withColumn("RECORD_FILE_NAME", col("_metadata.file_name"))
Pour obtenir le chemin d'accès complet du fichier :
.withColumn("RECORD_FILE_PATH", col("_metadata.file_path"))
Les deux options fonctionnent avec spark.read.
Les Notebook de Job effectuent des Opérations de données sur DBFS sur un cluster partagé
Lorsque vous utilisez DBFS avec un cluster partagé via le service FUSE, le cluster ne peut pas accéder au système de fichiers et génère une erreur de fichier introuvable.
Les exemples suivants échouent sur un cluster partagé :
with open('/dbfs/test/sample_file.csv', 'r') as file:
ls -ltr /dbfs/test
cat /dbfs/test/sample_file.csv
Utilisez l'une des solutions suivantes :
- Utilisez un volume Databricks Unity Catalog au lieu de DBFS (recommandé).
- Mettez à jour le code pour utiliser
dbutilsouspark, qui utilisent le chemin d'accès direct au stockage et bénéficient d'un accès à DBFS depuis les clusters partagés.