Se connecter à Lakebase
Aperçu
Cette fonctionnalité est en aperçu public.
Utilisez Structured Streaming pour écrire dans Lakebase 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 connecteur de sortie Lakebase pour des écritures en streaming à faible latence vers Lakebase. Ce récepteur ne vous oblige pas à implémenter des fonctions foreachBatch 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.
Pour synchroniser les données de Lakebase vers les tables Delta Lake dans le Lakehouse, dans la direction inverse, consultez Lakebase Change Data Feed.
Exigences
- Databricks Runtime 18 et versions ultérieures
- Compute classique avec modes d'accès dédié ou standard.
- Une base de données Lakebase
Se connecter à une base de données
Le sink Lakebase prend en charge les méthodes de connexion suivantes :
Tables Lakebase enregistrées avec Unity Catalog
Pour les tables Lakebase enregistrées avec Unity Catalog, le connecteur gère automatiquement les identifiants et utilise l'identité de l'utilisateur ou du Service Principal Databricks exécutant la query. Si la table n'existe pas, le connecteur crée la table.
Pour enregistrer une base de données Lakebase avec Unity Catalog, consultez Enregistrer une base de données Lakebase dans Unity Catalog.
Pour écrire dans une table Lakebase, utilisez la méthode .toTable() avec un nom de table entièrement qualifié, catalog.schema.table. L'exemple suivant montre les options requises, plus l'option facultative upsertkey :
- Python
- Scala
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
Remplacez les espaces réservés suivants :
<catalog>.<schema>.<table>: Le nom entièrement qualifié de la table cible. Lecatalogest le catalogue Unity Catalog que vous avez créé lors de l'enregistrement de la base de données Lakebase, consultez Enregistrer une base de données Lakebase dans Unity Catalog. Si la table n'existe pas, le connecteur la crée.<primary-key-column>Facultatif. Une liste de colonnes séparées par des virgules qui forment la clé upsert, par exempleidouuser_id,event_type. Si vous omettezupsertkey, la destination déduit la clé primaire de la table cible. Consultez le 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.
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 sur une table Lakebase, utilisez les options endpoint et dbtable. L'exemple suivant inclut également les options facultatives database et upsertkey :
- Python
- Scala
(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-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
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-column>") // Optional. Inferred from the table's primary key if omitted.
.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 leprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>format. Voir les identifiants compute.<database>Facultatif. Le nom de la base de données Postgres cible. La valeur par défaut estdatabricks_postgres. Consultez Gérer les bases de données.<schema>.<table>: La table cible au formatschema.table. Si vous omettez le schéma, le récepteur utilise le schémapublic. Utilisez des identifiants simples qui start par une lettre ou un trait de soulignement et ne contiennent que des lettres, des chiffres et des traits de soulignement ; les identifiants entre guillemets et les caractères spéciaux, tels que les tirets, ne sont pas pris en charge.<primary-key-column>Facultatif. Une liste de colonnes séparées par des virgules qui forment la clé upsert, par exempleidouuser_id,event_type. Si vous omettezupsertkey, la destination déduit la clé primaire de la table cible. Consultez le 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.
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 |
|---|---|---|
|
| Facultatif. Le délai maximal pour conserver les lignes dans le tampon avant le vidage. Par exemple : |
|
| Facultatif. Le nombre maximal de lignes pour chaque transaction de base de données. |
| Aucun | Obligatoire. Chemin d’accès à un répertoire de checkpoint, tel qu’un volume Unity Catalog ( |
| Aucun | Facultatif. Une liste de noms de colonnes séparés par des virgules qui constituent la clé d'upsert. Par exemple : |
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 |
|---|---|---|
|
| Facultatif. Nom de la base de données PostgreSQL cible. |
| Aucun | Obligatoire. Le nom de la table cible au format |
| Aucun | Obligatoire. L'Endpoint Lakebase, au format |
Comportement Upsert
Lorsque des clés d'insertion-mise à jour existent, qu'elles soient spécifiées avec upsertkey ou déduites par le récepteur à partir des clés primaires de la table, le récepteur effectue des opérations d'insertion-mise à jour dans la table avec la syntaxe INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ... de PostgreSQL.
Lorsqu'aucune clé d'upsert n'existe, le récepteur effectue des insertions. Le mode de sortie d'une query n'a aucun effet sur le comportement d'upsert ou d'insert.
Les upsertkey colonnes doivent :
- Être un sous-ensemble non vide des colonnes du DataFrame.
- Faites correspondre le
PRIMARY KEYde la table cible exactement. Si les colonnes que vous spécifiez ne correspondent pas à la clé primaire, la query échoue. - Doivent être des types comparables, tels que les types numériques ou chaîne. Pour éviter les interblocages de base de données lors des écritures concurrentes, le récepteur trie les lignes par clé d'upsert dans chaque batch. Les clés d'upsert ne prennent pas en charge les types complexes ou de structure.
Les noms de colonnes sont automatiquement mis entre guillemets doubles ", default dans PostgreSQL, ce qui permet de gérer les mots-clés réservés et les noms en casse mixte.
Les noms de table et de schéma doivent utiliser des identifiants simples qui start par une lettre ou un trait de soulignement et ne contenir que des lettres, des chiffres et des traits de soulignement. Le récepteur ne prend pas en charge les identifiants entre guillemets ou les caractères spéciaux, tels que les tirets, dans les noms de table ou de schéma.
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
batchsizelignes, qui est par default à1000. - L'âge du tampon dépasse
batchinterval, qui est100 millisecondsdefault.
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
batchintervalpour garantir un temps maximum plus court avant le flush. Voir Mode temps réel dans Structured Streaming pour les concepts et Exemples de mode temps réel pour un exemple de code. - Pour les charges de travail à haut throughput, augmentez
batchsizeafin 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 montre la prise en charge des types de Trigger Structured Streaming :
Déclencheur | Pris en charge |
|---|---|
| Oui |
| Oui |
| Oui |
| Oui |
Modes de sortie
Ce tableau indique la prise en charge des modes de sortie de Structured Streaming :
Mode de résultat | Pris en charge |
|---|---|
| Oui |
| Oui. Le comportement est identique à |
| Non |
Limitations
- Le compute Serverless et les LakeFlow Pipelines ne sont pas pris en charge.
- Seul Lakebase est pris en charge comme cible d'écriture. Les bases de données externes compatibles avec PostgreSQL ne sont pas prises en charge.