Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
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
- Databricks Runtime 18 y versiones posteriores
- Proceso clásico con modos de acceso dedicados o estándar.
- Una base de datos de Lakebase
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.cataloges 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 ejemploidouser_id,event_type. Si omiteupsertkey, 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 formatoprojects/<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 predeterminadodatabricks_postgres. Consulte Administración de bases de datos. -
<schema>.<table>: la tabla de destino en formatoschema.table. Si omite el esquema, el receptor usa el esquemapublic. 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 ejemploidouser_id,event_type. Si omiteupsertkey, 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 KEYcoincida 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
batchsizefilas, cuyo valor predeterminado es1000. - La antigüedad del búfer supera
batchinterval, cuyo valor predeterminado es100 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
batchintervalpara 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
batchsizepara 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.