Aller au contenu principal

Query de bases de données à l’aide de JDBC

Databricks prend en charge la connexion aux bases de données externes à l'aide de JDBC. Cet article fournit la syntaxe de base pour configurer et utiliser ces connexions avec des exemples en Python, SQL et Scala.

info

Expérimental

La documentation sur la fédération des query héritées a été retirée et pourrait ne pas être mise à jour. Les configurations mentionnées dans ce contenu ne sont ni officiellement approuvées ni testées par Databricks. Si la Lakehouse Federation prend en charge votre base de données source, Databricks recommande d'utiliser cette dernière.

Partner Connect fournit des intégrations optimisées pour synchroniser les données avec de nombreuses sources de données externes. Reportez-vous à Qu'est-ce que Databricks Partner Connect ?.

important

Les exemples de cet article n'incluent pas les noms d'utilisateur et les mots de passe dans les URL JDBC. Databricks recommande d'utiliser des secrets pour stocker vos identifiants de base de données. Par exemple :

Python
username = dbutils.secrets.get(scope = "jdbc", key = "username")
password = dbutils.secrets.get(scope = "jdbc", key = "password")

Pour référencer les secrets Databricks avec SQL, vous devez configurer une propriété de configuration Spark pendant l'initialisation du cluster.

Pour un exemple complet de gestion des secrets, consultez Tutoriel : Créer et utiliser un secret Databricks.

Établir la connectivité au cloud

Les Virtual Private Cloud (VPC) Databricks sont configurés pour autoriser uniquement les clusters Spark. Lorsque vous vous connectez à une autre infrastructure, la bonne pratique consiste à utiliser le peering de Virtual Private Cloud (VPC). Une fois le peering Virtual Private Cloud (VPC) établi, vous pouvez vérifier auprès du service public netcat sur le cluster.

Bash
%sh nc -vz <jdbcHostname> <jdbcPort>

Lire les données avec JDBC

Vous devez configurer un certain nombre de paramètres pour lire les données en utilisant JDBC. Notez que chaque base de données utilise un format différent pour le <jdbc-url>.

Python
employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)

Spark lit automatiquement le schéma de la table de base de données et mappe ses types aux types Spark SQL.

Python
employees_table.printSchema

Vous pouvez exécuter des queries sur cette table JDBC :

Python
display(employees_table.select("age", "salary").groupBy("age").avg("salary"))

Écrire des données avec JDBC

L'enregistrement de données dans des tables avec JDBC utilise des configurations similaires à la lecture. Voir l'exemple suivant :

Python
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.save()
)

Le comportement par default tente de créer une nouvelle table et génère une erreur si une table portant ce nom existe déjà.

Vous pouvez ajouter des données à une table existante à l'aide de la syntaxe suivante :

Python
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.mode("append")
.save()
)

Vous pouvez remplacer une table existante à l'aide de la syntaxe suivante :

Python
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.mode("overwrite")
.save()
)

Contrôler le parallélisme des queries JDBC

Par default, le Driver JDBC interroge la base de données source avec un seul thread. Pour améliorer les performances des lectures, vous devez spécifier un certain nombre d'options afin de contrôler le nombre de requêtes simultanées que Databricks envoie à votre base de données. Pour les petits clusters, la définition de l'option numPartitions égale au nombre de cœurs d'exécuteurs de votre cluster garantit que tous les nœuds query les données en parallèle.

attention

Définir numPartitions à une valeur élevée sur un grand cluster peut entraîner des performances négatives pour la base de données distante, car trop de queries simultanées pourraient submerger le service. Cela est particulièrement problématique pour les bases de données d'applications. Méfiez-vous de définir cette valeur au-dessus de 50.

remarque

Accélérez les query en sélectionnant une colonne avec un index calculé dans la base de données source pour le partitionColumn.

L'exemple de code suivant démontre la configuration du parallélisme pour un cluster à huit cœurs :

Python
employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
# a column that can be used that has a uniformly distributed range of values that can be used for parallelization
.option("partitionColumn", "<partition-key>")
# lowest value to pull data for with the partitionColumn
.option("lowerBound", "<min-value>")
# max value to pull data for with the partitionColumn
.option("upperBound", "<max-value>")
# number of partitions to distribute the data into. Do not set this very large (~hundreds)
.option("numPartitions", 8)
.load()
)
remarque

Databricks prend en charge toutes les options Apache Spark pour configurer JDBC.

Lors de l'écriture dans des bases de données à l'aide de JDBC, Apache Spark utilise le nombre de partitions en mémoire pour contrôler le parallélisme. Vous pouvez repartitionner les données avant l'écriture pour contrôler le parallélisme. Évitez un nombre élevé de partitions sur les grands clusters pour ne pas surcharger votre base de données distante. L'exemple suivant démontre le repartitionnement en huit partitions avant l'écriture :

Python
(employees_table.repartition(8)
.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.save()
)

Transférer une query vers le moteur de base de données

Vous pouvez déléguer une requête entière à la base de données et ne renvoyer que le résultat. Le paramètre table identifie la table JDBC à lire. Vous pouvez utiliser tout ce qui est valide dans une clause query SQL FROM.

Python
pushdown_query = "(select * from employees where emp_no < 10008) as emp_alias"

employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", pushdown_query)
.option("user", "<username>")
.option("password", "<password>")
.load()
)

Contrôlez le nombre de lignes récupérées par query

Les Drivers JDBC ont un paramètre fetchSize qui contrôle le nombre de lignes récupérées à la fois à partir de la base de données distante.

Paramètre

Résultat

Trop bas

Latence élevée due à de nombreux allers-retours (peu de lignes renvoyées par query)

Trop élevé

Erreur de mémoire (trop de données renvoyées dans une query)

Paramètre

Résultat

Trop bas

Latence élevée due à de nombreux allers-retours (peu de lignes renvoyées par query)

Trop élevé

Erreur de mémoire (trop de données renvoyées dans une query)

La valeur optimale dépend de la charge de travail. Les considérations comprennent :

  • Combien de colonnes sont renvoyées par la query ?
  • Quels types de données sont renvoyés ?
  • Quelle est la longueur des chaînes renvoyées dans chaque colonne ?

Les systèmes peuvent avoir des default très petits et bénéficier d'un ajustement. Par exemple : le default fetchSize d'Oracle est 10. Le fait de l'augmenter à 100 réduit le nombre total de query qui doivent être exécutées d'un facteur de 10. Les résultats JDBC sont du trafic réseau, il faut donc éviter les très grands nombres, mais les valeurs optimales pourraient être de l'ordre de milliers pour de nombreux datasets.

Utilisez l'option fetchSize, comme dans l'exemple suivant :

Python
employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.option("fetchSize", "100")
.load()
)