Aller au contenu principal

Se connecter à Lakebase

Utilisez Structured Streaming pour écrire dans Lakebase ou dans une base de données PostgreSQL externe avec le traitement par batch intégré, les nouvelles tentatives automatiques et l'authentification gérée par le workspace.

Quand utiliser le puits Lakebase

Utilisez le récepteur Lakebase pour des écritures en streaming à faible latence vers Lakebase ou une base de données PostgreSQL externe. Ce récepteur ne vous oblige pas à implémenter des fonctions foreach personnalisées pour gérer le traitement par batch, la gestion des connexions et la gestion des erreurs.

Quelques cas d’usage courants :

  • Mettez à jour les bases de données d'applications en temps réel pour les tableaux de bord opérationnels ou les fonctionnalités destinées aux clients.
  • Synchronisez les données qui changent continuellement, telles que les résultats de streaming agrégés ou filtrés, dans une base de données transactionnelle.
  • Écrivez le résultat d'une query Structured Streaming dans une table Lakebase avec une latence inférieure à la seconde en utilisant le mode temps réel.

Exigences

  • Databricks Runtime 18 LTS et versions supérieures.

    • Les connexions PostgreSQL externes et les types de données d’intervalle nécessitent Databricks Runtime 19 et versions ultérieures.
  • Compute classique avec des modes d'accès dédiés ou standard, ou compute serverless pour les notebooks ou les jobs. Sur le compute serverless, utilisez Trigger.AvailableNow(). Voir Streaming sur le compute serverless.

  • Une base de données Lakebase ou une connexion Unity Catalog à une base de données PostgreSQL externe.

Exigences relatives aux identifiants

Pour toutes les cibles, Databricks recommande d'utiliser des noms de schémas, de tables, de colonnes et de colonnes de clé primaire qui start par une lettre ou un trait de soulignement et ne contiennent que des lettres, des chiffres et des traits de soulignement. Le récepteur applique ces exigences lorsqu'il crée automatiquement une table Lakebase. Pour utiliser des identifiants ne respectant pas ces exigences, créez la table cible avant de start la query.

Se connecter à une base de données

Le sink Lakebase prend en charge les méthodes de connexion suivantes :

Tables Lakebase non enregistrées dans Unity Catalog

Pour les tables Lakebase non enregistrées auprès d'Unity Catalog, le connecteur gère automatiquement les identifiants et utilise l'identité de l'utilisateur ou du Service Principal Databricks exécutant la requête. Si la table n'existe pas, le connecteur crée la table.

Pour écrire dans une table Lakebase, utilisez les options endpoint et dbtable :

Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)

Remplacez les espaces réservés suivants :

  • <project-id>.<branch-id>.<endpoint-id>: votre Endpoint Lakebase. Trouvez les trois valeurs dans le Nom de la ressource du menu Obtenir l'ID de l'onglet Compute, qui a le projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id> format. Voir les identifiants compute.
  • <database>: facultatif. Le nom de la base de données PostgreSQL cible. Default to databricks_postgres. Voir Gérer les bases de données.
  • <schema>.<table>: La table cible au format schema.table. Si vous omettez le schéma, le récepteur utilise le schéma public. Pour la création automatique de tables, utilisez des identifiants qui start par une lettre ou un trait de soulignement et ne contiennent que des lettres, des chiffres et des traits de soulignement.
  • <primary-key-columns>: facultatif. Une liste séparée par des virgules de toutes les colonnes de la clé primaire de la table cible, par exemple id ou user_id,event_type. Voir Comportement d’upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: un chemin de volume Unity Catalog où la query stocke son point de contrôle. Vous pouvez également utiliser un URI de stockage d'objets cloud. L'emplacement doit être un stockage sur lequel vous pouvez écrire, pas un disque local, et doit être unique à chaque query de streaming. Ceci est indépendant de la table cible. Consultez la section Points de contrôle Structured Streaming.

Pour les configurations facultatives, telles que batchsize et batchinterval, consultez les options de configuration.

PostgreSQL externe avec des identifiants Unity Catalog

Dans Databricks Runtime 19 et versions supérieures, utilisez une connexion Unity Catalog pour vous authentifier auprès d'une base de données PostgreSQL externe sans stocker d'informations d'identification dans votre code. La table cible doit déjà exister.

Créer une connexion de type POSTGRESQL, voir Créer une connexion. L’utilisateur ou le Service Principal Databricks exécutant la query doit disposer de USE CONNECTION sur la connexion.

Pour écrire dans la table PostgreSQL, utilisez les options databricks.connection, database et dbtable :

Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)

Remplacez les espaces réservés suivants :

  • <connection-name>: le nom de la connexion Unity Catalog.
  • <database>: Le nom de la base de données PostgreSQL cible.
  • <schema>.<table>: La table cible existante au format schema.table. Si vous omettez le schéma, le récepteur utilise le schéma public.
  • <primary-key-columns>: facultatif. Une liste séparée par des virgules de toutes les colonnes de la clé primaire de la table cible, par exemple id ou user_id,event_type. Voir Comportement d’upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: un chemin de volume Unity Catalog où la query stocke son point de contrôle. Vous pouvez également utiliser un URI de stockage d'objets cloud. L'emplacement doit être un stockage sur lequel vous pouvez écrire, pas un disque local, et doit être unique à chaque query de streaming. Ceci est indépendant de la table cible. Consultez la section Points de contrôle Structured Streaming.

Les connexions PostgreSQL utilisent toujours TLS. La vérification du certificat suit les paramètres de la connexion Unity Catalog, que vous choisissez lorsque vous créez la connexion:

  • Certificat de serveur de confiance : lorsque cette option est sélectionnée, la connexion utilise sslmode=require, ce qui chiffre la connexion sans vérifier le certificat du serveur.
  • Certificat de serveur fourni par l’utilisateur : Fournissez un certificat de serveur encodé au format PEM à utiliser avec sslmode=verify-full lorsque Certificat de serveur de confiance n’est pas sélectionné. Si vous ne fournissez pas de certificat, la connexion utilise sslmode=verify-full avec le magasin de confiance par default de la JVM.

Options de configuration

Le récepteur génère une erreur pour les options non reconnues, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Les options suivantes s'appliquent à toutes les méthodes de connexion :

Clé

Par défaut

Description

batchinterval

100 milliseconds

Facultatif. Le délai maximal pour conserver les lignes dans le tampon avant le vidage. Par exemple : "50 milliseconds".

batchsize

1000

Facultatif. Le nombre maximal de lignes pour chaque transaction de base de données.

checkpointLocation

Aucun

Obligatoire. Chemin d’accès à un répertoire de checkpoint, tel qu’un volume Unity Catalog (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Doit être unique à chaque query. Consultez l’article dédié aux checkpoints Structured Streaming.

upsertkey

Aucun

Facultatif. Une liste séparée par des virgules de toutes les colonnes de la clé primaire de la table cible, par exemple "id" ou "user_id,event_type". Voir Comportement des upserts.

Clé

Par défaut

Description

batchinterval

100 milliseconds

Facultatif. Le délai maximal pour conserver les lignes dans le tampon avant le vidage. Par exemple : "50 milliseconds".

batchsize

1000

Facultatif. Le nombre maximal de lignes pour chaque transaction de base de données.

checkpointLocation

Aucun

Obligatoire. Chemin d’accès à un répertoire de checkpoint, tel qu’un volume Unity Catalog (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Doit être unique à chaque query. Consultez l’article dédié aux checkpoints Structured Streaming.

upsertkey

Aucun

Facultatif. Une liste séparée par des virgules de toutes les colonnes de la clé primaire de la table cible, par exemple "id" ou "user_id,event_type". Voir Comportement des upserts.

Tables Lakebase non enregistrées avec Unity Catalog

Les options suivantes s'appliquent lorsque vous vous connectez à une table Lakebase non enregistrée auprès de Unity Catalog:

Clé

Par défaut

Description

database

databricks_postgres

Facultatif. Nom de la base de données PostgreSQL cible.

dbtable

Aucun

Obligatoire. Le nom de la table cible au format schema.table. Si vous ne spécifiez pas de schéma, la valeur de schéma par default est public. Pour la création automatique de tables, utilisez des identifiants qui start par une lettre ou un trait de soulignement et contenant uniquement des lettres, des chiffres et des traits de soulignement.

endpoint

Aucun

Obligatoire. L'Endpoint Lakebase, au format project_id.branch_id ou project_id.branch_id.endpoint_id. Le paramètre endpoint_id est optionnel. Si vous l'omettez et que la Branch possède un unique Endpoint de lecture-écriture, le récepteur sélectionne cet Endpoint par default.

Clé

Par défaut

Description

database

databricks_postgres

Facultatif. Nom de la base de données PostgreSQL cible.

dbtable

Aucun

Obligatoire. Le nom de la table cible au format schema.table. Si vous ne spécifiez pas de schéma, la valeur de schéma par default est public. Pour la création automatique de tables, utilisez des identifiants qui start par une lettre ou un trait de soulignement et contenant uniquement des lettres, des chiffres et des traits de soulignement.

endpoint

Aucun

Obligatoire. L'Endpoint Lakebase, au format project_id.branch_id ou project_id.branch_id.endpoint_id. Le paramètre endpoint_id est optionnel. Si vous l'omettez et que la Branch possède un unique Endpoint de lecture-écriture, le récepteur sélectionne cet Endpoint par default.

PostgreSQL externe avec des identifiants Unity Catalog

Les options suivantes s’appliquent lorsque vous vous connectez à une base de données PostgreSQL externe avec des identifiants Unity Catalog:

Clé

Par défaut

Description

database

Aucun

Obligatoire. Le nom de la base de données PostgreSQL cible.

databricks.connection

Aucun

Obligatoire. Le nom de la connexion Unity Catalog pour l'authentification gérée par Unity Catalog vers PostgreSQL externe.

dbtable

Aucun

Obligatoire. Nom de la table cible existante au format schema.table. Si vous ne spécifiez pas de schéma, la valeur du schéma default est public.

Clé

Par défaut

Description

database

Aucun

Obligatoire. Le nom de la base de données PostgreSQL cible.

databricks.connection

Aucun

Obligatoire. Le nom de la connexion Unity Catalog pour l'authentification gérée par Unity Catalog vers PostgreSQL externe.

dbtable

Aucun

Obligatoire. Nom de la table cible existante au format schema.table. Si vous ne spécifiez pas de schéma, la valeur du schéma default est public.

Mappages de types de données

Le récepteur vérifie que chaque colonne DataFrame est compatible avec sa colonne cible correspondante avant d’écrire dans une table Lakebase ou PostgreSQL externe existante.

Le tableau suivant contient les types pris en charge dans Databricks Runtime 18 LTS et versions supérieures :

Type Spark

Type de table Lakebase créé automatiquement

Types compatibles dans les tables PostgreSQL existantes

ByteType, ShortType

smallint

smallint

IntegerType

integer

integer

LongType

bigint

bigint

FloatType

real

real

DoubleType

double precision

double precision

DecimalType

numeric

numeric

StringType

text

varchar, text

VarcharType(n)

varchar(n)

varchar, text

CharType(n)

char(n)

char

BinaryType

bytea

bytea

BooleanType

boolean

boolean

TimestampType

timestamptz

timestamptz

TimestampNTZType

timestamp

timestamp

DateType

date

date

ArrayType, MapType, StructType, VariantType, NullType

jsonb

json, jsonb

Type Spark

Type de table Lakebase créé automatiquement

Types compatibles dans les tables PostgreSQL existantes

ByteType, ShortType

smallint

smallint

IntegerType

integer

integer

LongType

bigint

bigint

FloatType

real

real

DoubleType

double precision

double precision

DecimalType

numeric

numeric

StringType

text

varchar, text

VarcharType(n)

varchar(n)

varchar, text

CharType(n)

char(n)

char

BinaryType

bytea

bytea

BooleanType

boolean

boolean

TimestampType

timestamptz

timestamptz

TimestampNTZType

timestamp

timestamp

DateType

date

date

ArrayType, MapType, StructType, VariantType, NullType

jsonb

json, jsonb

Le tableau suivant contient les types pris en charge dans Databricks Runtime 19 et versions supérieures :

Type Spark

Type de table Lakebase créé automatiquement

Types compatibles dans les tables PostgreSQL existantes

DayTimeIntervalType, YearMonthIntervalType

interval

interval

Type Spark

Type de table Lakebase créé automatiquement

Types compatibles dans les tables PostgreSQL existantes

DayTimeIntervalType, YearMonthIntervalType

interval

interval

Comportement d’upsert

L’option upsertkey identifie les colonnes de clé primaire de la table cible. Pour une table existante, les colonnes de upsertkey doivent correspondre exactement à la clé primaire de la table. Si vous omettez l’option, le récepteur lit la clé primaire à partir de la table. Pour une table Lakebase créée par le récepteur, upsertkey définit la clé primaire. Si vous omettez l’option, le récepteur crée la table sans clé primaire.

Lorsque la table cible possède une clé primaire, la destination effectue des upserts avec la syntaxe INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... de PostgreSQL. Lorsque la table cible ne possède aucune clé primaire, la destination effectue des insertions. Le mode de sortie d’une query n’a aucun effet sur ce comportement.

Toutes les colonnes de clé primaire doivent être présentes dans le DataFrame et utiliser des types comparables, tels que des types numériques ou de chaînes de caractères.

Optimisation des performances

Traitement par batch et contre-pression

Un vidage est Trigger lorsque l'une ou l'autre des conditions est remplie :

  • La mémoire tampon atteint batchsize lignes, qui est par default à 1000.
  • L'âge du tampon dépasse batchinterval, qui est 100 milliseconds default.

Lorsque la base de données ne peut pas suivre le débit de données entrant, le récepteur propage une rétropression en amont vers la source.

Conseils sur la latence et le throughput :

  • Pour les charges de travail à faible latence avec le mode temps réel, diminuez batchinterval afin de garantir un temps maximal plus court avant le vidage. Consultez Concepts du mode temps réel pour les concepts et Exemples du mode temps réel pour un exemple de code.
  • Pour les charges de travail à haut throughput, augmentez batchsize afin de réduire la surcharge pour chaque transaction.

Comportement de connexion

Le récepteur utilise le pool de connexions sur les exécuteurs. By default, chaque tâche utilise une connexion de base de données.

Databricks vous recommande d'utiliser la valeur par default de 1 tâche pour chaque connexion. Si vous augmentez le nombre de tâches pour chaque connexion, vous pourriez provoquer des contentions de connexion et augmenter les latences pour les connexions à haut throughput.

Pour configurer le rapport entre les tâches et les connexions, définissez la configuration Spark spark.databricks.sql.streaming.jdbc.tasksPerConnection. Si la base de données cible a une faible limite de connexions, réduisez le nombre de partitions de brassage ou augmentez spark.databricks.sql.streaming.jdbc.tasksPerConnection.

Le sink réessaie automatiquement les erreurs JDBC transitoires, y compris les échecs de connexion, les interblocages et la limitation du débit. Si le récepteur épuise toutes les tentatives, la query échoue.

Trigger et modes de sortie pris en charge

Trigger

Ce tableau indique la prise en charge des types de trigger Structured Streaming sur le compute classique et serverless :

Déclencheur

Classic Compute

Compute Serverless (notebooks et jobs)

RealTime

Oui

Non

ProcessingTime

Oui

Non

AvailableNow

Oui

Oui

Once

Oui. Obsolète. Utilisez AvailableNow.

Oui. Obsolète. Utilisez AvailableNow.

Déclencheur

Classic Compute

Compute Serverless (notebooks et jobs)

RealTime

Oui

Non

ProcessingTime

Oui

Non

AvailableNow

Oui

Oui

Once

Oui. Obsolète. Utilisez AvailableNow.

Oui. Obsolète. Utilisez AvailableNow.

Modes de sortie

Ce tableau indique la prise en charge des modes de sortie de Structured Streaming :

Mode de résultat

Pris en charge

update

Oui

append

Oui. Le comportement est identique à update. La query effectue un upsert lorsque la table cible possède une clé primaire ; sinon, la query réalise une insertion. Consultez Comportement de l’upsert.

complete

Non

Mode de résultat

Pris en charge

update

Oui

append

Oui. Le comportement est identique à update. La query effectue un upsert lorsque la table cible possède une clé primaire ; sinon, la query réalise une insertion. Consultez Comportement de l’upsert.

complete

Non

Limitations

  • Pour une base de données PostgreSQL externe connectée via une connexion Unity Catalog, la table cible doit déjà exister. Le récepteur crée automatiquement les tables manquantes uniquement dans Lakebase.
  • Les LakeFlow Pipelines ne sont pas pris en charge.