Utilisez un Notebook pour accéder à une instance de base de données
Lakebase Provisioned est l'offre originale de Lakebase qui utilise un compute provisionné que vous mettez à l'échelle manuellement. Pour les régions prises en charge, consultez la disponibilité des régions. Pour la dernière version de Lakebase, avec compute à dimensionnement automatique, mise à l'échelle jusqu'à zéro, création de branches et restauration instantanée, consultez Lakebase Autoscaling.
Depuis le 12 mars 2026, les nouvelles instances Lakebase sont créées en tant que projets de dimensionnement automatique. Les instances provisionnées existantes sont mises à niveau automatiquement vers la mise à l'échelle automatique, à compter de juin 2026. Pour plus de détails, consultez Mise à niveau vers le dimensionnement automatique Lakebase.
Cette page contient des exemples de code qui vous montrent comment accéder à votre instance de base de données Lakebase via les Notebooks Databricks et exécuter des queries à l'aide de Python et Scala.
Les exemples couvrent différentes stratégies de connexion pour s'adapter à différents cas d'utilisation :
- Connexion unique : utilisée pour les scripts simples où une connexion unique à la base de données est ouverte, utilisée et fermée.
- Connection pool : utilisé pour les charges de travail à high concurrency, où un Pool de connexions réutilisables est maintenu.
- Rotation du jeton OAuth M2M : utilise des jetons OAuth de courte durée, automatiquement actualisés pour l'authentification.
Les exemples suivants génèrent des identifiants sécurisés par programmation. Évitez d'inclure directement des identifiants dans un notebook. Databricks vous recommande d'utiliser l'une des méthodes sécurisées suivantes :
- Stockez les mots de passe Postgres dans les secrets Databricks.
- Générez des jetons OAuth en utilisant M2M OAuth.
Avant de commencer
Assurez-vous de satisfaire aux exigences suivantes avant d'accéder à votre instance de base de données :
- Vous disposez d'un rôle Postgres correspondant pour vous connecter à l'instance de base de données. Consultez les rôles Postgres.
- Votre rôle Postgres se voit accorder les autorisations nécessaires pour accéder à la base de données, au schéma ou à la table.
- Vous pouvez vous authentifier à l'instance de la base de données. Si vous devez obtenir manuellement un jeton OAuth pour votre instance de base de données, consultez Authentification à une instance de base de données.
Si vous utilisez Private Link, vous devez utiliser un cluster mono-utilisateur.
Python
Le SDK Python Databricks peut être utilisé pour obtenir un jeton OAuth pour une instance de base de données respective.
Connectez-vous à votre instance de base de données à partir d'un Notebook Databricks à l'aide des bibliothèques Python suivantes :
psycopg(psycopg3)psycopgavec pooling de connexionsSQLAlchemy
Prérequis
Avant d'exécuter les exemples de code suivants, mettez à niveau le SDK Databricks pour Python vers la version 0,61,0 ou supérieure, puis redémarrez Python.
%pip install databricks-sdk>=0.61.0
%restart_python
psycopg
Les exemples de code illustrent une seule connexion et l'utilisation d'un Pool de connexions. Pour en savoir plus sur la façon d'obtenir l'instance de base de données et les informations d'identification par programme, consultez la façon d'obtenir un jeton OAuth à l'aide du SDK Python.
%pip install "psycopg[binary,pool]"
- Single connection
- Connection pool
import psycopg
from databricks.sdk import WorkspaceClient
import uuid
w = WorkspaceClient()
instance_name = "<YOUR INSTANCE>"
instance = w.database.get_database_instance(name=instance_name)
cred = w.database.generate_database_credential(request_id=str(uuid.uuid4()), instance_names=[instance_name])
# Connection parameters
conn = psycopg.connect(
host = instance.read_write_dns,
dbname = "databricks_postgres",
user = "<YOUR USER>",
password = cred.token,
sslmode = "require"
)
# Execute query
with conn.cursor() as cur:
cur.execute("SELECT version()")
version = cur.fetchone()[0]
print(version)
conn.close()
import psycopg
from psycopg_pool import ConnectionPool
from databricks.sdk import WorkspaceClient
import uuid
w = WorkspaceClient()
instance_name = "<YOUR INSTANCE>"
instance = w.database.get_database_instance(name=instance_name)
cred = w.database.generate_database_credential(request_id=str(uuid.uuid4()), instance_names=[instance_name])
# Create a connection pool
connection_pool = ConnectionPool(
conninfo=f"dbname=databricks_postgres user=<YOUR USER> host={instance.read_write_dns} port=5432 password={cred.token} sslmode=require",
min_size=1, # Minimum number of connections in the pool
max_size=10, # Maximum number of connections in the pool
open=True
)
print("Connection pool created successfully")
def executeWithPgConnection(execFn):
with connection_pool.connection() as connection:
print("Successfully received a connection from the pool")
execFn(connection)
print("Connection returned to the pool")
def printVersion(connection):
cursor = connection.cursor()
cursor.execute("SELECT version()")
version = cursor.fetchone()
print(f"Connected to PostgreSQL database. Version: {version}")
executeWithPgConnection(printVersion)
psycopg3
L'exemple de code illustre l'utilisation d'un pool de connexions avec un OAuth M2M rotatif. Il utilise generate_database_credential(). Pour en savoir plus sur la façon d'obtenir l'instance de base de données et les identifiants par programmation, consultez comment obtenir un jeton OAuth à l'aide du SDK Python.
%pip install "psycopg[binary,pool]"
from databricks.sdk import WorkspaceClient
import uuid
import psycopg
import string
from psycopg_pool import ConnectionPool
w = WorkspaceClient()
class CustomConnection(psycopg.Connection):
global w
def __init__(self, *args, **kwargs):
# Call the parent class constructor
super().__init__(*args, **kwargs)
@classmethod
def connect(cls, conninfo='', **kwargs):
# Append the new password to kwargs
cred = w.database.generate_database_credential(request_id=str(uuid.uuid4()), instance_names=[instance_name])
kwargs['password'] = cred.token
# Call the superclass's connect method with updated kwargs
return super().connect(conninfo, **kwargs)
username = "<YOUR USER>"
instance_name = "<YOUR INSTANCE>"
instance = w.database.get_database_instance(name=instance_name)
host = instance.read_write_dns
port = 5432
database = "databricks_postgres"
pool = ConnectionPool(
conninfo=f"dbname={database} user={username} host={host}",
connection_class=CustomConnection,
min_size=1,
max_size=10,
open=True
)
with pool.connection() as conn:
with conn.cursor() as cursor:
cursor.execute("SELECT version()")
for record in cursor:
print(record)
SQLAlchemy
Les exemples de code démontrent une connexion unique et l'utilisation d'un pool de connexions avec un jeton OAuth M2M rotatif. Pour en savoir plus sur la façon d'obtenir l'instance de base de données et les informations d'identification par programme, consultez la façon d'obtenir un jeton OAuth à l'aide du SDK Python.
- Single connection
- Connection pool & rotating M2M OAuth
%pip install "sqlalchemy>=2.0" "psycopg[binary]"
from sqlalchemy import create_engine, text
from databricks.sdk import WorkspaceClient
import uuid
w = WorkspaceClient()
instance_name = "<YOUR INSTANCE>"
instance = w.database.get_database_instance(name=instance_name)
cred = w.database.generate_database_credential(request_id=str(uuid.uuid4()), instance_names=[instance_name])
user = "<YOUR USER>"
host = instance.read_write_dns
port = 5432
database = "databricks_postgres"
password = cred.token
connection_pool = create_engine(f"postgresql+psycopg://{user}:{password}@{host}:{port}/{database}?sslmode=require")
with connection_pool.connect() as conn:
result = conn.execute(text("SELECT version()"))
for row in result:
print(f"Connected to PostgreSQL database. Version: {row}")
%pip install "sqlalchemy>=2.0" "psycopg[binary]"
from databricks.sdk import WorkspaceClient
import uuid
import time
from sqlalchemy import create_engine, text, event
w = WorkspaceClient()
instance_name = "<YOUR INSTANCE>"
instance = w.database.get_database_instance(name=instance_name)
username = "<YOUR USER>"
host = instance.read_write_dns
port = 5432
database = "databricks_postgres"
# sqlalchemy setup + function to refresh the OAuth token that is used as the Postgres password every 15 minutes.
connection_pool = create_engine(f"postgresql+psycopg://{username}:@{host}:{port}/{database}")
postgres_password = None
last_password_refresh = time.time()
@event.listens_for(connection_pool, "do_connect")
def provide_token(dialect, conn_rec, cargs, cparams):
global postgres_password, last_password_refresh, host
if postgres_password is None or time.time() - last_password_refresh > 900:
print("Refreshing PostgreSQL OAuth token")
cred = w.database.generate_database_credential(request_id=str(uuid.uuid4()), instance_names=[instance_name])
postgres_password = cred.token
last_password_refresh = time.time()
cparams["password"] = postgres_password
with connection_pool.connect() as conn:
result = conn.execute(text("SELECT version()"))
for row in result:
print(f"Connected to PostgreSQL database. Version: {row}")
Scala
Les exemples de code montrent comment obtenir par programmation l'instance de base de données et les identifiants, et comment se connecter à une instance de base de données à l'aide d'une connexion unique ou d'un Pool de connexions.
Étape 1 : Utilisez le SDK Java Databricks pour obtenir un jeton OAuth
Pour plus de détails sur la façon d'obtenir l'instance de base de données et les informations d'identification par programmation, consultez comment obtenir un jeton OAuth à l'aide du SDK Java.
Étape 2 : se connecter à une instance de base de données
- Single connection
- Connection pool
import java.sql.{Connection, DriverManager, ResultSet, Statement}
Class.forName("org.postgresql.Driver")
val user = "<YOUR USER>"
val host = instance.getName()
val port = "5432"
val database = "databricks_postgres"
val password = cred.getToken()
val url = f"jdbc:postgresql://${host}:${port}/${database}"
val connection = DriverManager.getConnection(url, user, password)
println("Connected to PostgreSQL database!")
val statement = connection.createStatement()
val resultSet = statement.executeQuery("SELECT version()")
if (resultSet.next()) {
val version = resultSet.getString(1)
println(s"PostgreSQL version: $version")
}
import com.zaxxer.hikari.{HikariConfig, HikariDataSource}
import java.sql.Connection
// Configure HikariCP
val config = new HikariConfig()
config.setJdbcUrl("jdbc:postgresql://instance.getName():5432/databricks_postgres")
config.setUsername("<YOUR USER>")
config.setPassword(cred.getToken())
config.setMaximumPoolSize(10)
// Create a data source
val dataSource = new HikariDataSource(config)
// Function to get a connection and execute a query
def runQuery(): Unit = {
var connection: Connection = null
try {
// Get a connection from the pool
connection = dataSource.getConnection()
// Create a statement
val statement = connection.createStatement()
// Execute a query
val resultSet = statement.executeQuery("SELECT version() AS v;")
// Process the result set
while (resultSet.next()) {
val v = resultSet.getString("v")
println(s"*******Connected to PostgreSQL database. Version: $v")
}
} catch {
case e: Exception => e.printStackTrace()
} finally {
// Close the connection which returns it to the pool
if (connection != null) connection.close()
}
}
// Run the query
runQuery()
// Close the data source
dataSource.close()