Se connecter à Lakebase

Utilisez le Structured Streaming pour écrire sur Lakebase ou une base de données PostgreSQL externe avec batching intégré, essais automatiques et authentification gérée par l’espace de travail.

Quand utiliser le récepteur Lakebase

Utilisez le puits 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 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

  • Databricks Runtime 18 LTS et plus.
    • Les connexions PostgreSQL externes nécessitent l’utilisation de Databricks Runtime 19 ou version ultérieure et de s’inscrire à la préversion Custom JDBC on UC Compute.
    • Les types de données à intervalles exigent d’utiliser Databricks Runtime 19 et supérieur.
  • Calcul classique avec modes d’accès dédiés ou standards, ou calcul sans serveur pour ordinateurs portables ou tâches. Sur le calcul sans serveur, utilisez Trigger.AvailableNow(). Consultez Streaming sur l’informatique sans serveur.
  • Une base de données Lakebase, ou une connexion Unity Catalog à une base de données PostgreSQL externe.

Exigences d’identification

Pour toutes les cibles, Databricks recommande d’utiliser des noms de schéma, tableau, colonne et colonne de clé primaire qui commencent par une lettre ou un sous-trait et ne contiennent que des lettres, des chiffres et des sous-traits. Le sink applique ces exigences lorsqu’il crée automatiquement une table Lakebase. Pour utiliser des identifiants qui ne répondent pas à ces exigences, créez la table cible avant de commencer la requête.

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 sur une table Lakebase, utilisez la .toTable() méthode avec un nom de table entièrement qualifié, catalog.schema.table:

Python

(df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")
)

Scala

df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .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-columns>: facultatif. Une liste séparée par 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 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 optionnelles, telles que batchsize et batchinterval, voir les options du récepteur PostgreSQL.

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

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-columns>")  // Optional
  .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. Le nom de la base de données cible PostgreSQL. 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. Pour la création automatique de tables, utilisez des identifiants qui commencent 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 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 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 des configurations facultatives, telles que batchsize et batchinterval, voir les options du récepteur PostgreSQL.

PostgreSQL externe avec identifiants du catalogue Unity

Important

Cette fonctionnalité est disponible en préversion publique. Les administrateurs d’espace de travail peuvent contrôler l’accès à Custom JDBC sur UC Compute depuis la page Preview . Consultez Gérer les préversions d’Azure Databricks.

Utilisez une connexion Unity Catalog pour vous authentifier à une base de données PostgreSQL externe sans stocker d’identifiants 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 principal de service qui exécute la requête doit disposer de USE CONNECTION sur la connexion.

Pour écrire dans la table PostgreSQL, utilisez les databricks.connectionoptions , 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()
)

Scala

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 placeholders 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 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 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.

Les connexions PostgreSQL utilisent toujours TLS. La vérification des certificats suit les paramètres de la connexion Unity Catalog, que vous choisissez lors de la création de la connexion :

  • Certificat serveur de confiance : Une fois sélectionné, la connexion utilise sslmode=require, qui chiffre la connexion sans vérifier le certificat serveur.
  • Certificat serveur fourni par l’utilisateur : Fournir un certificat serveur encodé en PEM à utiliser sslmode=verify-full lorsque le certificat 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 défaut de la JVM.

Options de configuration

Le récepteur signale une erreur pour les options non reconnues, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Pour les options de configuration du sink, y compris les options communes et celles pour chaque méthode de connexion, voir les options du sink PostgreSQL.

Mappages de types de données

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

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

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

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 clés principales de la table cible. Pour une table existante, les colonnes dans upsertkey doivent correspondre exactement à la clé primaire de la table. Si vous omettez cette option, le sink lit la clé primaire de la table. Pour une table Lakebase créée par le récepteur, upsertkey définit la clé primaire. Si vous omettez cette option, le récepteur crée la table sans définir de clé primaire.

Lorsque la table cible possède une clé primaire, le récepteur effectue une opération d’upsert avec la syntaxe INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... de PostgreSQL. Lorsque la table cible n’a pas de clé primaire, le puits effectue des insertions. Le mode de sortie d’une requête n’a aucun effet sur ce comportement.

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

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. Voir concepts en mode temps réel pour les concepts et exemples de mode temps réel pour 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 indique la prise en charge des types de déclencheurs de Structured Streaming sur les calculs classiques et sans serveur :

Déclencheur Calcul classique Calcul sans serveur (carnets et tâches)
RealTime Yes No
ProcessingTime Yes No
AvailableNow Yes Yes
Once Yes. Deprecated. Utilisez AvailableNow. Yes. Deprecated. Utilisez AvailableNow.

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. Voir comportement d’Upsert.
complete No

Limites

  • Pour une base de données PostgreSQL externe connectée via une connexion Unity Catalog, la table cible doit déjà exister. L’évier crée automatiquement des tables manquantes uniquement dans Lakebase.
  • Les pipelines Lakeflow ne sont pas pris en charge.