Se connecter à Lakebase

Important

Cette fonctionnalité est disponible en préversion publique.

Utilisez Structured Streaming pour écrire dans Lakebase avec le traitement par lots intégré, les nouvelles tentatives automatiques et l’authentification gérée par l’espace de travail.

Quand utiliser le récepteur 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 lots, les connexions et le traitement des erreurs.

Les cas d’utilisation courants sont les suivants :

  • Mettez à jour les bases de données d’application en temps réel pour les tableaux de bord opérationnels ou les fonctionnalités orientées client.
  • Synchronisez la modification continue des données, telles que les résultats de diffusion en continu agrégés ou filtrés, dans une base de données transactionnelle.
  • Écrivez la sortie d’une requête Structured Streaming dans une table Lakebase avec une latence de sous-seconde à l’aide du mode en temps réel.

Pour synchroniser des données de Lakebase vers des tables Delta Lake dans le Lakehouse, pour la direction inverse, voir Lakebase Change Data Feed.

Exigences

Se connecter à une base de données

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

Tables Lakebase inscrites auprès du catalogue Unity

Pour les tables Lakebase inscrites auprès du catalogue Unity, le connecteur gère automatiquement les informations d’identification et utilise l’identité de l’utilisateur ou du principal du service exécutant la requête. Si la table n’existe pas, le connecteur crée la table.

Pour inscrire une base de données Lakebase auprès du catalogue Unity, consultez Inscrire une base de données Lakebase dans le catalogue Unity.

Pour écrire dans une table Lakebase, utilisez la méthode .toTable() avec un nom de table pleinement qualifié, catalog.schema.table. L’exemple suivant montre les options requises, ainsi que 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>")
)

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>")

Remplacez les placeholders suivants :

  • <catalog>.<schema>.<table> : Le nom pleinement qualifié de la table cible. catalog correspond au catalogue Unity Catalog que vous avez créé lorsque vous avez enregistré la base de données Lakebase ; voir 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. Liste séparée par des virgules des colonnes qui forment la clé upsert, par exemple id ou user_id,event_type. Si vous omettez upsertkey, le récepteur déduit la clé de la clé primaire de la table cible. Consultez le comportement d’Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Un chemin d’accès vers un volume d’Unity Catalog dans lequel la requête stocke son checkpoint. Vous pouvez également utiliser un URI de stockage d’objets cloud. L’emplacement doit être le stockage dans lequel vous pouvez écrire, et non sur un disque local, et doit être unique à chaque requête de diffusion en continu. Ceci est indépendant de la table cible. Consultez les Points de contrôle Structured Streaming.

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

Tables Lakebase non inscrites auprès du catalogue Unity

Pour les tables Lakebase non inscrites auprès du catalogue Unity, le connecteur gère automatiquement les informations d’identification et utilise l’identité de l’utilisateur ou du principal du service 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. 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()
)

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()

Remplacez les placeholders suivants :

  • <project-id>.<branch-id>.<endpoint-id> : votre point de terminaison Lakebase. Recherchez les trois valeurs dans le nom de la ressource dans le menu Obtenir l’ID de l’onglet Calculs , qui a le format projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Consultez les identifiants de calcul.
  • <database>: facultatif. 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>: tableau cible au schema.table format. Si vous omettez le schéma, le récepteur utilise le schéma public. Utilisez des identificateurs simples qui commencent par une lettre ou un trait de soulignement et contiennent uniquement des lettres, des chiffres et des traits de soulignement ; Les identificateurs entre guillemets et les caractères spéciaux, tels que les traits d’union, ne sont pas pris en charge.
  • <primary-key-column>: facultatif. Liste séparée par des virgules des colonnes qui forment la clé upsert, par exemple id ou user_id,event_type. Si vous omettez upsertkey, le récepteur déduit la clé de la clé primaire de la table cible. Consultez le comportement d’Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Un chemin d’accès vers un volume d’Unity Catalog dans lequel la requête stocke son checkpoint. Vous pouvez également utiliser un URI de stockage d’objets cloud. L’emplacement doit être le stockage dans lequel vous pouvez écrire, et non sur un disque local, et doit être unique à chaque requête de diffusion en continu. Ceci est indépendant de la table cible. Consultez les 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 signale une erreur pour les options non reconnues, JDBC_STREAMING_SINK_INVALID_OPTIONS.

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

Clé Default Description
batchinterval 100 milliseconds Optional. Durée maximale pendant laquelle les lignes sont conservées dans la mémoire tampon avant vidage. Par exemple : "50 milliseconds".
batchsize 1000 Optional. Nombre maximal de lignes pour chaque transaction de base de données.
checkpointLocation None Required. Chemin d’accès à un répertoire de point de contrôle, tel qu’un volume de catalogue Unity (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Doit être unique à chaque requête. Consultez les Points de contrôle Structured Streaming.
upsertkey None Optional. Liste séparée par des virgules des noms de colonnes qui forment la clé upsert. Par exemple, "id" ou "user_id,event_type". Si vous spécifiez upsertkey, les colonnes doivent correspondre à la clé primaire de la table, ou la requête échoue. Si vous l’omettez, le récepteur utilise automatiquement la clé primaire. Pour plus d’informations, consultez le comportement Upsert.

Tables Lakebase non inscrites auprès du catalogue Unity

Les options suivantes s’appliquent lorsque vous vous connectez à une table Lakebase non inscrite auprès du catalogue Unity :

Clé Default Description
database databricks_postgres Optional. Nom de la base de données PostgreSQL cible.
dbtable None Required. Nom de la table cible au format schema.table. Si vous ne spécifiez pas de schéma, la valeur de schéma par défaut est public. Utilisez des identificateurs simples qui commencent par une lettre ou un trait de soulignement et contiennent uniquement des lettres, des chiffres et des traits de soulignement. Ne mettez pas entre guillemets les noms de table ou de schéma ; les identificateurs entre guillemets et les noms contenant des caractères spéciaux, tels que des traits d’union, ne sont pas pris en charge.
endpoint None Required. Le point de terminaison Lakebase, au format project_id.branch_id ou project_id.branch_id.endpoint_id. Le endpoint_id est facultatif ; si vous l’omettez et que la branche possède un seul point de terminaison en lecture-écriture, le collecteur sélectionne ce point de terminaison par défaut.

Comportement d’upsert

Lorsque des clés d’upsert existent, qu’elles soient spécifiées avec upsertkey ou déduites par le connecteur de destination à partir des clés primaires de la table, le connecteur de destination effectue des opérations d’upsert dans la table à l’aide de 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 requête n’a aucun effet sur le comportement d’upsert ou d’insertion.

Les upsertkey colonnes doivent :

  • Être un sous-ensemble non vide des colonnes DataFrame.
  • Correspond exactement à la PRIMARY KEY de la table cible. Si les colonnes que vous spécifiez ne correspondent pas à la clé primaire, la requête échoue.
  • Être des types comparables, tels que des types numériques ou de chaînes. Pour empêcher les interblocages de base de données lors des écritures simultanées, le récepteur trie les lignes par clé d’upsert au sein de chaque lot. Les clés d’upsert ne prennent pas en charge pas les types complexes ni les types struct.

Les noms de colonnes sont automatiquement mis entre guillemets doubles ", conformément au comportement par défaut de PostgreSQL, ce qui permet de gérer les mots-clés réservés et les noms mêlant majuscules et minuscules.

Les noms de table et de schéma doivent utiliser des identificateurs simples qui commencent par une lettre ou un trait de soulignement et contiennent uniquement des lettres, des chiffres et des traits de soulignement. Le récepteur ne prend pas en charge les identifiants entre guillemets ni les caractères spéciaux, tels que les tirets, dans les noms de tables ou de schémas.

Réglage des performances

Traitement par lots et contre-pression

Un vidage est déclenché lorsque l’une des deux conditions est satisfaite :

  • La mémoire tampon contient jusqu’à batchsize lignes, avec une valeur par défaut de 1000.
  • L’âge de la mémoire tampon dépasse batchinterval, ce qui est défini par défaut sur 100 milliseconds.

Lorsque la base de données ne peut pas suivre le débit de données entrant, le collecteur propage une contre-pression en amont vers la source.

Conseils sur la latence et le débit :

  • Pour les charges de travail à faible latence avec le mode en temps réel, diminuez batchinterval pour garantir une durée maximale plus courte avant le vidage. Consultez le mode temps réel dans Structured Streaming pour connaître les concepts et les exemples de mode temps réel d’un exemple de code.
  • Pour les charges de travail à haut débit, augmentez batchsize pour réduire la surcharge liée à chaque transaction.

Comportement de connexion

Le collecteur utilise un pool de connexions au niveau des exécuteurs. Par défaut, chaque tâche utilise une connexion de base de données.

Databricks recommande d’utiliser la valeur de tâche par défaut 1 pour chaque connexion. Si vous augmentez le nombre de tâches pour chaque connexion, vous pouvez provoquer des conflits de connexion et augmenter les latences pour les connexions à haut débit.

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

Le récepteur effectue automatiquement des tentatives de reprise pour les erreurs JDBC transitoires, notamment les échecs de connexion, les interblocages et les limitations du débit. Si le récepteur épuise toutes les tentatives, la requête échoue.

Déclencheurs et modes de sortie pris en charge

Triggers

Ce tableau présente la prise en charge des types de déclencheurs Structured Streaming :

Déclencheur Soutenu
realTime Yes
ProcessingTime Yes
AvailableNow Yes
Once Yes

Modes de sortie

Ce tableau présente la prise en charge des modes de sortie Structured Streaming :

Mode de sortie Soutenu
update Yes
append Yes. Le comportement est identique à update. La requête effectue des upserts si la table cible possède une clé primaire, sinon elle effectue des insertions. Consultez le comportement d’Upsert.
complete No

Limites

  • Le calcul serverless et les pipelines Lakeflow ne sont pas pris en charge.
  • Seul Lakebase est pris en charge en tant que cible d’écriture. Les bases de données compatibles Avec PostgreSQL externes ne sont pas prises en charge.