Aller au contenu principal

Se connecter à Lakebase

info

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
(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. Le catalog est 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 exemple id ou user_id,event_type. Si vous omettez upsertkey, 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
(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 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 Postgres cible. La valeur par défaut est databricks_postgres. Consultez 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. 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 exemple id ou user_id,event_type. Si vous omettez upsertkey, 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

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 de noms de colonnes séparés par des virgules qui constituent la clé d'upsert. Par exemple : "id" ou "user_id,event_type". Si vous spécifiez upsertkey, les colonnes doivent correspondre à la clé primaire de la table, sinon la query échoue. Si vous l'omettez, le récepteur utilise automatiquement la clé primaire. Pour plus d'information, consultez le comportement d'upsert.

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 de noms de colonnes séparés par des virgules qui constituent la clé d'upsert. Par exemple : "id" ou "user_id,event_type". Si vous spécifiez upsertkey, les colonnes doivent correspondre à la clé primaire de la table, sinon la query échoue. Si vous l'omettez, le récepteur utilise automatiquement la clé primaire. Pour plus d'information, consultez le comportement d'upsert.

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. 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. Ne mettez pas les noms de table ou de schéma entre guillemets ; les identifiants entre guillemets et les noms avec des caractères spéciaux, tels que les tirets, ne sont pas pris en charge.

endpoint

Aucun

Obligatoire. L'Endpoint Lakebase, au format project_id.branch_id ou project_id.branch_id.endpoint_id. L'élément endpoint_id est facultatif ; si vous l'omettez et que la Branch dispose d'un Endpoint en lecture-écriture unique, le récepteur sélectionne cet Endpoint by 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. 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. Ne mettez pas les noms de table ou de schéma entre guillemets ; les identifiants entre guillemets et les noms avec des caractères spéciaux, tels que les tirets, ne sont pas pris en charge.

endpoint

Aucun

Obligatoire. L'Endpoint Lakebase, au format project_id.branch_id ou project_id.branch_id.endpoint_id. L'élément endpoint_id est facultatif ; si vous l'omettez et que la Branch dispose d'un Endpoint en lecture-écriture unique, le récepteur sélectionne cet Endpoint by default.

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 KEY de 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 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 pour 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 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 montre la prise en charge des types de Trigger Structured Streaming :

Déclencheur

Pris en charge

realTime

Oui

ProcessingTime

Oui

AvailableNow

Oui

Once

Oui

Déclencheur

Pris en charge

realTime

Oui

ProcessingTime

Oui

AvailableNow

Oui

Once

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

update

Oui

append

Oui. Le comportement est identique à update. La query effectue des upserts lorsque la table cible a une clé primaire ; sinon, la query insère. Consultez le comportement d'Upsert.

complete

Non

Mode de résultat

Pris en charge

update

Oui

append

Oui. Le comportement est identique à update. La query effectue des upserts lorsque la table cible a une clé primaire ; sinon, la query insère. Consultez le comportement d'Upsert.

complete

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.