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.
Utiliza Structured Streaming para escribir en Lakebase o en una base de datos PostgreSQL externa con procesamiento por lotes integrado, intentos automáticos y autenticación gestionada por el espacio de trabajo.
Cuándo usar el receptor Lakebase
Usa el sumidero de Lakebase para escrituras de streaming de baja latencia a Lakebase o a una base de datos externa de PostgreSQL. Este receptor no requiere que implemente funciones foreach 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 LTS y superiores.
- Las conexiones externas de PostgreSQL requieren que uses Databricks Runtime 19 y superiores y que optes por la vista previa de JDBC personalizada en UC Compute .
- Los tipos de datos de intervalo requieren que uses Databricks Runtime 19 y superiores.
- Computación clásica con modos de acceso dedicados o estándar, o computación sin servidor para cuadernos o tareas. En computación sin servidor, usa
Trigger.AvailableNow(). Consulte Streaming en computación sin servidor. - Una base de datos Lakebase, o una conexión de Unity Catalog a una base de datos PostgreSQL externa.
Requisitos de identificador
Para todos los objetivos, Databricks recomienda usar nombres de esquema, tabla, columna y columnas de clave primaria que comiencen por una letra o guion bajo y contengan solo letras, números y guiones bajos. El sumidero hace cumplir estos requisitos cuando crea automáticamente una tabla Lakebase. Para usar identificadores que no cumplan estos requisitos, crea la tabla objetivo antes de iniciar la consulta.
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 Lakebase, utiliza el .toTable() método con un nombre de tabla completamente cualificado, 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>")
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-columns>: opcional. Una lista separada por comas de todas las columnas en la clave primaria de la tabla objetivo, por ejemploidouser_id,event_type. Ver 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, utiliza las opciones endpoint y 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()
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. El nombre de la base de datos objetivo PostgreSQL. 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. Para la creación automática de tablas, utiliza identificadores que empiecen por una letra o guion bajo y contengan solo letras, números y guiones bajos. -
<primary-key-columns>: opcional. Una lista separada por comas de todas las columnas en la clave primaria de la tabla objetivo, por ejemploidouser_id,event_type. Ver 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.
PostgreSQL externo con credenciales del Catálogo de Unity
Importante
Esta característica está en versión preliminar pública. Los administradores de espacios de trabajo pueden controlar el acceso a JDBC personalizado en UC Compute desde la página de Previsualizaciones . Consulte Administrar versiones preliminares de Azure Databricks.
Usa una conexión de Unity Catalog para autenticarte en una base de datos PostgreSQL externa sin almacenar credenciales en tu código. La tabla objetivo debe existir ya.
Crear una conexión de tipo POSTGRESQL, ver Crear una conexión. El usuario o la entidad de servicio que ejecuta la consulta debe tener USE CONNECTION en la conexión.
Para escribir en la tabla de PostgreSQL, utilice las opciones databricks.connection, database y 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()
Sustituya los siguientes marcadores de posición:
-
<connection-name>: El nombre de la conexión del Unity Catalog. -
<database>: El nombre de la base de datos objetivo PostgreSQL. -
<schema>.<table>: La tabla de destino existente en formatoschema.table. Si omite el esquema, el receptor usa el esquemapublic. -
<primary-key-columns>: opcional. Una lista separada por comas de todas las columnas en la clave primaria de la tabla objetivo, por ejemploidouser_id,event_type. Ver 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.
Las conexiones PostgreSQL siempre usan TLS. La verificación de certificados sigue la configuración de la conexión del Catálogo de Unity, que eliges al crear la conexión:
-
Certificado del servidor de confianza: Cuando se selecciona, la conexión utiliza
sslmode=require, que cifra la conexión sin verificar el certificado del servidor. -
Certificado de servidor proporcionado por el usuario: Proporciona un certificado de servidor codificado en PEM para usar
sslmode=verify-fullcuando no se seleccione el certificado del servidor de confianza . Si no proporcionas un certificado, la conexión usasslmode=verify-fullcon el almacén de confianza predeterminado de la JVM.
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 separada por comas de todas las columnas en la clave primaria de la tabla objetivo, por ejemplo, "id" o "user_id,event_type". Ver 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. Para la creación automática de tablas, utiliza identificadores que empiecen por una letra o guion bajo y contengan solo letras, números y guiones bajos. |
endpoint |
Ninguno | Required. El punto de conexión de Lakebase, en formato project_id.branch_id o project_id.branch_id.endpoint_id.
endpoint_id es opcional. Si lo omites y la rama tiene un único extremo de lectura-escritura, el sumidero selecciona ese extremo por defecto. |
PostgreSQL externo con credenciales del Catálogo de Unity
Las siguientes opciones se aplican cuando se conecta a una base de datos externa de PostgreSQL con credenciales del Unity Catalog:
| Clave | Predeterminado | Description |
|---|---|---|
database |
Ninguno | Required. Nombre de la base de datos postgreSQL de destino. |
databricks.connection |
Ninguno | Required. El nombre de conexión del Catálogo de Unity para la autenticación gestionada por el Catálogo de Unity a PostgreSQL externo. |
dbtable |
Ninguno | Required. El nombre de la tabla de destino existente con formato schema.table. Si no especifica un esquema, el valor de esquema predeterminado es public. |
Mapeos de tipos de datos
El sumidero comprueba que cada columna DataFrame sea compatible con su columna objetivo correspondiente antes de escribir en una tabla existente de Lakebase o PostgreSQL externa.
La siguiente tabla contiene los tipos soportados en Databricks Runtime 18 LTS y superiores:
| Tipo de Spark | Tipo de tabla Lakebase creado automáticamente | Tipos compatibles en tablas PostgreSQL existentes |
|---|---|---|
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 |
La siguiente tabla contiene los tipos soportados en Databricks Runtime 19 y superiores:
| Tipo de Spark | Tipo de tabla Lakebase creado automáticamente | Tipos compatibles en tablas PostgreSQL existentes |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Comportamiento de Upsert
La upsertkey opción identifica las columnas clave principales de la tabla objetivo. Para una tabla existente, las columnas en upsertkey deben coincidir exactamente con la clave primaria de la tabla. Si omites la opción, el fregadero lee la clave primaria de la tabla. Para una tabla Lakebase que crea el sumidero, upsertkey define la clave primaria. Si omites la opción, el sink crea la tabla sin una clave primaria.
Cuando la tabla de destino tiene una clave primaria, el destino realiza operaciones upsert mediante la sintaxis INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... de PostgreSQL. Cuando la tabla destino no tiene clave primaria, el sumidero realiza insertos. El modo de salida de una consulta no afecta a este comportamiento.
Todas las columnas clave primarias deben estar presentes en el DataFrame y usar tipos comparables, como tipos numéricos o de cadena.
Ajuste de 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. Consulta conceptos de modo en tiempo real para conceptos y ejemplos de modo en tiempo real para 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 desencadenador de Structured Streaming en la computación clásica y sin servidor:
| Desencadenador | Computación clásica | Computación sin servidor (portátiles y trabajos) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Yes. Deprecated. Utilice AvailableNow. |
Yes. Deprecated. Utilice AvailableNow. |
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. Ver comportamiento de Upsert. |
complete |
No |
Limitaciones
- Para una base de datos PostgreSQL externa conectada a través de una conexión al Unity Catalog, la tabla objetivo debe existir ya. El fregadero crea automáticamente tablas faltantes solo en Lakebase.
- Las canalizaciones de Lakeflow no son compatibles.