Conexión a Lakebase

Importante

Esta característica está en versión preliminar pública.

Use Structured Streaming para escribir en Lakebase con procesamiento por lotes integrado, reintentos automáticos y autenticación administrada por el área de trabajo.

Cuándo usar el receptor Lakebase

Utilice el sumidero de Lakebase para realizar escrituras en streaming de baja latencia en Lakebase. Este receptor no requiere que implemente funciones foreachBatch personalizadas para controlar el procesamiento por lotes, la administración de conexiones y el control de errores.

Entre los casos de uso comunes se incluyen:

  • Actualice las bases de datos de aplicaciones en tiempo real para los paneles operativos o las características orientadas al cliente.
  • Sincronice los datos que cambian continuamente, como los resultados de streaming agregados o filtrados, en una base de datos transaccional.
  • Escriba el resultado de una consulta de Structured Streaming en una tabla de Lakebase con una latencia inferior a un segundo utilizando el modo en tiempo real.

Para sincronizar datos desde Lakebase a tablas de Delta Lake en Lakehouse, consulte fuente de cambios de datos de Lakebase para la dirección inversa.

Requisitos

Conectar con una base de datos

El sumidero de Lakebase admite los siguientes métodos de conexión:

Tablas de Lakebase registradas con el catálogo de Unity

En el caso de las tablas de Lakebase registradas con el catálogo de Unity, el conector administra automáticamente las credenciales y usa la identidad del usuario o la entidad de servicio que ejecuta la consulta. Si la tabla no existe, el conector crea la tabla.

Para registrar una base de datos de Lakebase con el catálogo de Unity, consulte Registro de una base de datos de Lakebase en el catálogo de Unity.

Para escribir en una tabla de Lakebase, utilice el método .toTable() con un nombre de tabla completo, catalog.schema.table. En el ejemplo siguiente se muestran las opciones necesarias, además de la opción opcional 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>")

Sustituya los siguientes marcadores de posición:

  • <catalog>.<schema>.<table>: nombre completo de la tabla de destino. catalog es el catálogo de Unity Catalog que creó al registrar la base de datos de Lakebase; consulte Registrar una base de datos de Lakebase en Unity Catalog. Si la tabla no existe, el conector lo crea.
  • <primary-key-column>: opcional. Lista separada por comas de las columnas que forman la clave upsert, por ejemplo id o user_id,event_type. Si omite upsertkey, el receptor infiere la clave a partir de la clave primaria de la tabla de destino. Consulte Comportamiento de Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: ruta de un volumen de Unity Catalog donde la consulta almacena su punto de control. También puede usar un URI de almacenamiento de objetos en la nube. La ubicación debe ser el almacenamiento en el que puede escribir, no en el disco local y debe ser único para cada consulta de streaming. Esto es independiente de la tabla de destino. Consulte Puntos de control de Structured Streaming.

Para obtener configuraciones opcionales, como batchsize y batchinterval, vea Opciones de configuración.

Tablas de Lakebase no registradas con el catálogo de Unity

En el caso de las tablas de Lakebase no registradas con el catálogo de Unity, el conector administra automáticamente las credenciales y usa la identidad del usuario o la entidad de servicio que ejecuta la consulta. Si la tabla no existe, el conector crea la tabla.

Para escribir en una tabla de Lakebase, utilice las opciones endpoint y dbtable. En el ejemplo siguiente también se incluyen las opciones opcionales database y 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()

Sustituya los siguientes marcadores de posición:

  • <project-id>.<branch-id>.<endpoint-id>: Tu punto de conexión de Lakebase. Busque los tres valores en el nombre del recurso en el menú Obtener identificador de la pestaña Proceso , que tiene el formato projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Consulte Identificadores de proceso.
  • <database>: opcional. Nombre de la base de datos postgres de destino. Tiene como valor predeterminado databricks_postgres. Consulte Administración de bases de datos.
  • <schema>.<table>: la tabla de destino en formato schema.table. Si omite el esquema, el receptor usa el esquema public. Use identificadores simples que empiecen por una letra o un carácter de subrayado y contengan solo letras, números y caracteres de subrayado; No se admiten identificadores entre comillas y caracteres especiales, como guiones.
  • <primary-key-column>: opcional. Lista separada por comas de las columnas que forman la clave upsert, por ejemplo id o user_id,event_type. Si omite upsertkey, el receptor infiere la clave a partir de la clave primaria de la tabla de destino. Consulte Comportamiento de Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: ruta de un volumen de Unity Catalog donde la consulta almacena su punto de control. También puede usar un URI de almacenamiento de objetos en la nube. La ubicación debe ser el almacenamiento en el que puede escribir, no en el disco local y debe ser único para cada consulta de streaming. Esto es independiente de la tabla de destino. Consulte Puntos de control de Structured Streaming.

Para obtener configuraciones opcionales, como batchsize y batchinterval, vea Opciones de configuración.

Opciones de configuración

El sumidero genera un error si hay opciones no reconocidas, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Las siguientes opciones se aplican a todos los métodos de conexión:

Clave Predeterminado Description
batchinterval 100 milliseconds Optional. Tiempo máximo para contener filas en el búfer antes de vaciar. Por ejemplo: "50 milliseconds".
batchsize 1000 Optional. Número máximo de filas para cada transacción de base de datos.
checkpointLocation Ninguno Required. Ruta de acceso a un directorio de punto de control, como un volumen de catálogo de Unity (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Debe ser único para cada consulta. Consulte Puntos de control de Structured Streaming.
upsertkey Ninguno Optional. Una lista de nombres de las columnas, separados por comas, que forman la clave de upsert. Por ejemplo, "id" o "user_id,event_type". Si especifica upsertkey, las columnas deben coincidir con la clave principal de la tabla o se produce un error en la consulta. Si se omite, el receptor usa automáticamente la clave principal. Para obtener más información, consulte Comportamiento de Upsert.

Tablas de Lakebase no registradas con el catálogo de Unity

Las siguientes opciones se aplican al conectarse a una tabla de Lakebase no registrada con el catálogo de Unity:

Clave Predeterminado Description
database databricks_postgres Optional. Nombre de la base de datos postgreSQL de destino.
dbtable Ninguno Required. Nombre de la tabla de destino en schema.table formato. Si no especifica un esquema, el valor de esquema predeterminado es public. Use identificadores simples que empiecen por una letra o un carácter de subrayado y contengan solo letras, números y caracteres de subrayado. No utilice comillas en los nombres de tablas o esquemas; los identificadores entre comillas y los nombres con caracteres especiales, como los guiones, no son compatibles.
endpoint Ninguno Required. El punto de conexión de Lakebase, en formato project_id.branch_id o project_id.branch_id.endpoint_id. El endpoint_id es opcional; si se omite y la rama tiene un único punto de conexión de lectura y escritura, el receptor selecciona ese punto de conexión de forma predeterminada.

Comportamiento de Upsert

Cuando existen claves de upsert, ya sea especificadas con upsertkey o inferidas por el receptor a partir de las claves primarias de la tabla, el receptor realiza operaciones de upsert en la tabla con la sintaxis INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ... de PostgreSQL.

Cuando no existen claves upsert, el receptor realiza inserciones. El modo de salida de una consulta no tiene ningún efecto en el comportamiento upsert o insert.

Las upsertkey columnas deben:

  • Ser un subconjunto no vacío de las columnas DataFrame.
  • Haz que la tabla de destino PRIMARY KEY coincida exactamente. Si las columnas especificadas no coinciden con la clave principal, se produce un error en la consulta.
  • Ser tipos comparables, como tipos numéricos o de cadena. Para evitar interbloqueos en la base de datos durante las escrituras concurrentes, el receptor ordena las filas por clave de actualización o inserción (upsert) dentro de cada lote. Las claves Upsert no admiten tipos complejos o de estructura.

Los nombres de columna se entrecomillan automáticamente con el valor predeterminado de PostgreSQL, las comillas dobles ", lo que permite manejar palabras clave reservadas y nombres con combinación de mayúsculas y minúsculas.

Los nombres de tabla y esquema deben usar identificadores simples que empiecen por una letra o un carácter de subrayado y solo contengan letras, números y caracteres de subrayado. El receptor no admite identificadores entre comillas ni caracteres especiales, como guiones, en los nombres de tabla o de esquema.

Optimización del rendimiento

Procesamiento por lotes y contrapresión

Se desencadena un vaciamiento cuando se cumple cualquiera de las condiciones:

  • El búfer alcanza las batchsize filas, cuyo valor predeterminado es 1000.
  • La antigüedad del búfer supera batchinterval, cuyo valor predeterminado es 100 milliseconds.

Cuando la base de datos no puede seguir el ritmo de la tasa de entrada de datos, el sumidero propaga contrapresión aguas arriba hacia la fuente.

Guía de latencia y rendimiento:

  • Para cargas de trabajo de baja latencia con modo en tiempo real, disminuya batchinterval para garantizar un tiempo máximo más corto antes del vaciado. Consulte Modo en tiempo real en Structured Streaming para obtener conceptos y ejemplos de modo en tiempo real para obtener un ejemplo de código.
  • Para cargas de trabajo de alto rendimiento, aumente batchsize para reducir la sobrecarga de cada transacción.

Comportamiento de la conexión

El receptor usa la agrupación de conexiones en los ejecutores. De forma predeterminada, cada tarea usa una conexión de base de datos.

Databricks recomienda que utilice el valor predeterminado de 1 para cada tarea de conexión. Si aumenta el número de tareas de cada conexión, puede provocar contenciones de conexión y aumentar las latencias de las conexiones de alto rendimiento.

Para configurar la relación de tareas con las conexiones, establezca la spark.databricks.sql.streaming.jdbc.tasksPerConnection configuración de Spark. Si la base de datos de destino tiene un límite bajo de conexiones, reduzca el número de particiones de mezcla o aumente spark.databricks.sql.streaming.jdbc.tasksPerConnection.

El receptor reintenta automáticamente los errores transitorios de JDBC, incluidos los errores de conexión, los interbloqueos y la limitación de velocidad. Si el receptor agota todos los reintentos, la consulta falla.

Desencadenadores y modos de salida admitidos

Triggers

Esta tabla muestra la compatibilidad con los tipos de activación de Structured Streaming:

Desencadenador Compatible
realTime Yes
ProcessingTime Yes
AvailableNow Yes
Once Yes

Modos de salida

En esta tabla se muestra la compatibilidad con los modos de salida de Structured Streaming:

Modo de salida Compatible
update Yes
append Yes. El comportamiento es idéntico a update. La consulta realiza una operación de actualización o inserción cuando la tabla de destino tiene una clave primaria; de lo contrario, inserta. Consulte Comportamiento de Upsert.
complete No

Limitaciones

  • La computación sin servidor y las canalizaciones de Lakeflow no son compatibles.
  • Solo Lakebase es compatible como destino de escritura. No se admiten bases de datos externas compatibles con PostgreSQL.