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.
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 ?.
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
- Scala
username = dbutils.secrets.get(scope = "jdbc", key = "username")
password = dbutils.secrets.get(scope = "jdbc", key = "password")
val username = dbutils.secrets.get(scope = "jdbc", key = "username")
val 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.
%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
- SQL
- Scala
employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)
CREATE TEMPORARY VIEW employees_table_vw
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'
)
val 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
- SQL
- Scala
employees_table.printSchema
DESCRIBE employees_table_vw
employees_table.printSchema
Vous pouvez exécuter des queries sur cette table JDBC :
- Python
- SQL
- Scala
display(employees_table.select("age", "salary").groupBy("age").avg("salary"))
SELECT age, avg(salary) as salary
FROM employees_table_vw
GROUP BY age
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
- SQL
- Scala
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.save()
)
CREATE TABLE new_employees_table
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'
) AS
SELECT * FROM employees_table_vw
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
- SQL
- Scala
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.mode("append")
.save()
)
CREATE TABLE IF NOT EXISTS new_employees_table
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'
);
INSERT INTO new_employees_table
SELECT * FROM employees_table_vw;
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
- SQL
- Scala
(employees_table.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.mode("overwrite")
.save()
)
CREATE OR REPLACE TABLE new_employees_table
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'
) AS
SELECT * FROM employees_table_vw;
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.
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.
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
- SQL
- Scala
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()
)
CREATE TEMPORARY VIEW employees_table_vw
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>',
partitionColumn "<partition-key>",
lowerBound "<min-value>",
upperBound "<max-value>",
numPartitions 8
)
val 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()
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
- SQL
- Scala
(employees_table.repartition(8)
.write
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<new-table-name>")
.option("user", "<username>")
.option("password", "<password>")
.save()
)
CREATE TABLE new_employees_table
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'
) AS
SELECT /*+ REPARTITION(8) */ * FROM employees_table_vw
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
- SQL
- Scala
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()
)
CREATE TEMPORARY VIEW employees_table_vw
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "(select * from employees where emp_no < 10008) as emp_alias",
user '<username>',
password '<password>'
)
val pushdown_query = "(select * from employees where emp_no < 10008) as emp_alias"
val 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) |
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
- SQL
- Scala
employees_table = (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.option("fetchSize", "100")
.load()
)
CREATE TEMPORARY VIEW employees_table_vw
USING JDBC
OPTIONS (
url "<jdbc-url>",
dbtable "<table-name>",
user '<username>',
password '<password>'.
fetchSize 100
)
val employees_table = spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.option("fetchSize", "100")
.load()