Conexión a Lakebase

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. 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-columns>: opcional. 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.
  • /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 formato projects/<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 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. 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 ejemplo id o user_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 formato schema.table. Si omite el esquema, el receptor usa el esquema public.
  • <primary-key-columns>: opcional. 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.
  • /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-full cuando no se seleccione el certificado del servidor de confianza . Si no proporcionas un certificado, la conexión usa sslmode=verify-full con 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 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. 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 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 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.