Pular para o conteúdo principal

Conectar ao Lakebase

Use o Structured Streaming para gravar no Lakebase ou em um banco de dados PostgreSQL externo com processamento em lotes integrado, novas tentativas automáticas e autenticação gerenciada pelo workspace.

Quando usar o destino do Lakebase

Use o coletor do Lakebase para gravações em transmissão de baixa latência no Lakebase ou em um banco de dados PostgreSQL externo. Este coletor não exige que você implemente funções customizadas foreach para gerenciar processamento em lote, gerenciamento de conexões e tratamento de erros.

Casos de uso comuns incluem:

  • Atualize bancos de dados de aplicativos em tempo real para painéis operacionais ou recursos voltados para o cliente.
  • Sincronizar dados em constante alteração, como resultados de transmissão agregados ou filtrados, em um banco de dados transacional.
  • Grave a saída de uma consulta de transmissão estructurada em uma tabela Lakebase com latência de subsegundo usando modo em tempo real.

Para sincronizar dados do Lakebase com as tabelas do Delta Lake no Lakehouse, no sentido inverso, consulte Lakebase Change Data Feed.

Requisitos

  • Databricks Runtime 18 LTS ou acima.

  • Compute clássico com modos de acesso dedicados ou padrão, ou compute serverless para notebooks ou jobs. Em compute serverless, use Trigger.AvailableNow(). Consulte Transmissão em compute serverless.

  • Um banco de dados Lakebase ou uma conexão do Unity Catalog com um banco de dados PostgreSQL externo.

Requisitos de identificador

Para todos os destinos, o Databricks recomenda o uso de nomes de esquema, tabela, coluna e coluna de chave primária que comecem com uma letra ou underscore e contenham apenas letras, números e underscores. O coletor impõe esses requisitos quando cria automaticamente uma tabela Lakebase. Para usar identificadores que não atendem a esses requisitos, crie a tabela de destino antes de começar a query.

Conecte-se a um banco de dados

O coletor do Lakebase oferece suporte aos seguintes métodos de conexão:

Tabelas do Lakebase registradas no Unity Catalog

Para tabelas Lakebase registradas no Unity Catalog, o conector gerencia automaticamente as credenciais e usa a identidade do usuário ou da entidade de serviço do Databricks executando a consulta. Se a tabela não existir, o conector cria a tabela.

Para registrar um banco de dados Lakebase no Unity Catalog, consulte Registrar um banco de dados Lakebase no Unity Catalog.

Para gravar em uma tabela do Lakebase, use o método .toTable() com um nome de tabela totalmente qualificado, 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>")
)

Substitua os seguintes espaços reservados:

  • <catalog>.<schema>.<table>: O nome totalmente qualificado da tabela de destino. O catalog é o catálogo do Unity Catalog que você criou quando registrou o banco de dados Lakebase. Consulte Registrar um banco de dados Lakebase no Unity Catalog. Se a tabela não existir, o conector a cria.
  • <primary-key-columns>: Opcional. Uma lista separada por vírgulas de todas as colunas na chave principal da tabela de destino, por exemplo, id ou user_id,event_type. Consulte Comportamento de upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho de volume do Unity Catalog onde a consulta armazena seu ponto de verificação. Você também pode usar um URI de armazenamento de objeto em cloud. O local deve ser um armazenamento onde você pode gravar, não um disco local, e deve ser exclusivo para cada consulta de transmissão. Isso é independente da tabela de destino. Consulte pontos de verificação de transmissão estructurada.

Para configurações opcionais, como batchsize e batchinterval, consulte Opções de configuração.

Tabelas do Lakebase não registradas no Unity Catalog

Para tabelas Lakebase não registradas no Unity Catalog, o conector gerencia automaticamente as credenciais e usa a identidade do usuário ou da entidade de serviço Databricks que está executando a consulta. Se a tabela não existe, o conector cria a tabela.

Para gravar em uma tabela do Lakebase, use as opções endpoint e 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()
)

Substitua os seguintes espaços reservados:

  • <project-id>.<branch-id>.<endpoint-id>: Seu endpoint Lakebase. Encontre todos os três valores no nome do recurso no menu Obter ID da tab Computes , que tem o formato projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Consulte Identificadores de Compute.
  • <database>: Opcional. O nome do banco de dados PostgreSQL de destino. O default é databricks_postgres. Consulte Gerenciar bancos de dados.
  • <schema>.<table>: A tabela de destino no formato schema.table. Se você omitir o esquema, o sink usará o esquema public. Para a criação automática de tabelas, use identificadores que comecem com uma letra ou underscore e contenham apenas letras, números e underscores.
  • <primary-key-columns>: Opcional. Uma lista separada por vírgulas de todas as colunas na chave principal da tabela de destino, por exemplo, id ou user_id,event_type. Consulte Comportamento de upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho de volume do Unity Catalog onde a consulta armazena seu ponto de verificação. Você também pode usar um URI de armazenamento de objeto em cloud. O local deve ser um armazenamento onde você pode gravar, não um disco local, e deve ser exclusivo para cada consulta de transmissão. Isso é independente da tabela de destino. Consulte pontos de verificação de transmissão estructurada.

Para configurações opcionais, como batchsize e batchinterval, consulte Opções de configuração.

External PostgreSQL with Unity Catalog credentials

No Databricks Runtime 19 e superior, use uma conexão do Unity Catalog para se autenticar em um banco de dados PostgreSQL externo sem armazenar credenciais em seu código. A tabela de destino já deve existir.

Crie uma conexão do tipo POSTGRESQL, consulte Criar uma conexão. O usuário ou o Service Principal do Databricks que executa a query deve ter USE CONNECTION na conexão.

To write to the PostgreSQL table, use the databricks.connection, database, and dbtable options:

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

Substitua os seguintes espaços reservados:

  • <connection-name>: O nome da conexão do Unity Catalog.
  • <database>: O nome do banco de dados PostgreSQL de destino.
  • <schema>.<table>: A tabela de destino existente no formato schema.table. Se você omitir o esquema, o sink usará o esquema public.
  • <primary-key-columns>: Opcional. Uma lista separada por vírgulas de todas as colunas na chave principal da tabela de destino, por exemplo, id ou user_id,event_type. Consulte Comportamento de upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho de volume do Unity Catalog onde a consulta armazena seu ponto de verificação. Você também pode usar um URI de armazenamento de objeto em cloud. O local deve ser um armazenamento onde você pode gravar, não um disco local, e deve ser exclusivo para cada consulta de transmissão. Isso é independente da tabela de destino. Consulte pontos de verificação de transmissão estructurada.

As conexões PostgreSQL sempre usam TLS. A verificação de certificado segue as configurações da conexão do Unity Catalog, que você escolhe ao criar a conexão:

  • Certificado de servidor confiável : quando selecionada, a conexão usa sslmode=require, que criptografa a conexão sem verificar o certificado do servidor.
  • Certificado de servidor fornecido pelo usuário : forneça um certificado de servidor codificado em PEM para usar sslmode=verify-full quando a opção Certificado de servidor confiável não estiver selecionada. Se você não fornecer um certificado, a conexão usará sslmode=verify-full com o armazenamento de confiança default da JVM.

Opções de configuração

O sink gera um erro para opções não reconhecidas, JDBC_STREAMING_SINK_INVALID_OPTIONS.

As seguintes opções se aplicam a todos os métodos de conexão:

Chave

Padrão

Descrição

batchinterval

100 milliseconds

Opcional. O tempo máximo para reter linhas no buffer antes de liberar. Por exemplo, "50 milliseconds".

batchsize

1000

Opcional. O número máximo de linhas para cada transação de banco de dados.

checkpointLocation

Nenhuma

Obrigatório. Caminho para um diretório de ponto de verificação, como um volume do Unity Catalog (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Deve ser exclusivo para cada consulta. Consulte Pontos de verificação de Transmissão estructurada.

upsertkey

Nenhuma

Opcional. Uma lista separada por vírgulas de todas as colunas na chave primária da tabela de destino, por exemplo, "id" ou "user_id,event_type". Consulte Comportamento de upsert.

Chave

Padrão

Descrição

batchinterval

100 milliseconds

Opcional. O tempo máximo para reter linhas no buffer antes de liberar. Por exemplo, "50 milliseconds".

batchsize

1000

Opcional. O número máximo de linhas para cada transação de banco de dados.

checkpointLocation

Nenhuma

Obrigatório. Caminho para um diretório de ponto de verificação, como um volume do Unity Catalog (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Deve ser exclusivo para cada consulta. Consulte Pontos de verificação de Transmissão estructurada.

upsertkey

Nenhuma

Opcional. Uma lista separada por vírgulas de todas as colunas na chave primária da tabela de destino, por exemplo, "id" ou "user_id,event_type". Consulte Comportamento de upsert.

Tabelas do Lakebase não registradas com o Unity Catalog

As seguintes opções se aplicam ao se conectar a uma tabela do Lakebase não registrada com o Unity Catalog:

Chave

Padrão

Descrição

database

databricks_postgres

Opcional. O nome do banco de dados PostgreSQL de destino.

dbtable

Nenhuma

Obrigatório. O nome da tabela de destino no formato schema.table. Se você não especificar um esquema, o valor do esquema default será public. Para a criação automática de tabelas, use identificadores que comecem com uma letra ou underscore e contenham apenas letras, números e underscores.

endpoint

Nenhuma

Obrigatório. O endpoint do Lakebase, no formato project_id.branch_id ou project_id.branch_id.endpoint_id. O endpoint_id é opcional. Se você o omitir e a branch tiver um único endpoint de leitura e gravação, o coletor selecionará esse endpoint por default.

Chave

Padrão

Descrição

database

databricks_postgres

Opcional. O nome do banco de dados PostgreSQL de destino.

dbtable

Nenhuma

Obrigatório. O nome da tabela de destino no formato schema.table. Se você não especificar um esquema, o valor do esquema default será public. Para a criação automática de tabelas, use identificadores que comecem com uma letra ou underscore e contenham apenas letras, números e underscores.

endpoint

Nenhuma

Obrigatório. O endpoint do Lakebase, no formato project_id.branch_id ou project_id.branch_id.endpoint_id. O endpoint_id é opcional. Se você o omitir e a branch tiver um único endpoint de leitura e gravação, o coletor selecionará esse endpoint por default.

PostgreSQL externo com credenciais do Unity Catalog

As seguintes opções se aplicam quando você se conecta a um banco de dados PostgreSQL externo com credenciais do Unity Catalog:

Chave

Padrão

Descrição

database

Nenhuma

Obrigatório. O nome do banco de dados PostgreSQL de destino.

databricks.connection

Nenhuma

Obrigatório. O nome da conexão do Unity Catalog para autenticação gerenciada pelo Unity Catalog no PostgreSQL externo.

dbtable

Nenhuma

Obrigatório. O nome da tabela de destino existente no formato schema.table. Se você não especificar um esquema, o valor do esquema default será public.

Chave

Padrão

Descrição

database

Nenhuma

Obrigatório. O nome do banco de dados PostgreSQL de destino.

databricks.connection

Nenhuma

Obrigatório. O nome da conexão do Unity Catalog para autenticação gerenciada pelo Unity Catalog no PostgreSQL externo.

dbtable

Nenhuma

Obrigatório. O nome da tabela de destino existente no formato schema.table. Se você não especificar um esquema, o valor do esquema default será public.

Mapeamentos de tipos de dados

O coletor verifica se cada coluna do DataFrame é compatível com sua coluna de destino correspondente antes de gravar em um Lakebase existente ou em um banco de dados PostgreSQL externo.

A tabela a seguir contém os tipos com suporte no Databricks Runtime 18 LTS e acima:

Tipo Spark

Automatically created Lakebase table type

Compatible types in existing PostgreSQL tables

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

Tipo Spark

Automatically created Lakebase table type

Compatible types in existing PostgreSQL tables

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

A tabela a seguir contém os tipos com suporte no Databricks Runtime 19 e acima:

Tipo Spark

Automatically created Lakebase table type

Compatible types in existing PostgreSQL tables

DayTimeIntervalType, YearMonthIntervalType

interval

interval

Tipo Spark

Automatically created Lakebase table type

Compatible types in existing PostgreSQL tables

DayTimeIntervalType, YearMonthIntervalType

interval

interval

Comportamento de upsert

A opção upsertkey identifica as colunas da chave principal da tabela de destino. Para uma tabela existente, as colunas em upsertkey devem corresponder exatamente à chave principal da tabela. Se você omitir a opção, o coletor (sink) lerá a chave principal da tabela. Para uma tabela Lakebase criada pelo coletor, upsertkey define a chave principal. Se você omitir a opção, o coletor cria a tabela sem uma chave principal.

Quando a tabela de destino tem uma chave primária, o coletor realiza upsert com a sintaxe INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... do PostgreSQL. Quando a tabela de destino não tem uma chave primária, o coletor executa inserções. O modo de saída de uma query não tem efeito sobre esse comportamento.

Todas as colunas de primary key devem estar presentes no DataFrame e usar tipos comparáveis, como tipos numéricos ou de strings.

Desempenho tuning

Processamento em lote e contrapressão

Um flush é acionado quando uma das condições é atendida:

  • O buffer atinge batchsize linhas, cujo padrão é 1000.
  • A idade do buffer excede batchinterval, que por default é 100 milliseconds.

Quando o banco de dados não consegue acompanhar a taxa de dados de entrada, o destino propaga contrapressão a montante para a fonte.

Latência e Taxa de transferência: Orientações

  • Para cargas de trabalho de baixa latência com modo em tempo real, diminua batchinterval para garantir um tempo máximo menor antes da liberação. Consulte Conceitos do modo em tempo real para obter conceitos e Exemplos do modo em tempo real para obter um exemplo de código.
  • Para cargas de trabalho de alta taxa de transferência, aumente batchsize para reduzir a sobrecarga de cada transação.

Comportamento da conexão

O sink usa pooling de conexão em executors. Por padrão, cada tarefa usa uma conexão de banco de dados.

O Databricks recomenda que você use o valor default de 1 tarefa para cada conexão. Ao aumentar o número de tarefas para cada conexão, poderá causar contenções de conexão e aumentar as latências para conexões de alta taxa de transferência.

Para configurar a proporção de tarefas para conexões, defina a configuração do Spark spark.databricks.sql.streaming.jdbc.tasksPerConnection. Se o banco de dados de destino tiver um limite de conexão baixo, reduza o número de partições em ordem aleatória ou aumente spark.databricks.sql.streaming.jdbc.tasksPerConnection.

O sink tenta novamente, de forma automática, erros transitórios de JDBC, incluindo falhas de conexão, deadlocks e limitação de taxa. Se o coletor esgotar todas as tentativas, a consulta falhará.

Gatilhos e modos de saída compatíveis

Gatilhos

Esta tabela mostra o suporte para tipos de Trigger do Structured Streaming em compute classic e serverless:

Trigger

Compute clássico

Compute serverless (notebooks e jobs)

RealTime

Sim

Não

ProcessingTime

Sim

Não

AvailableNow

Sim

Sim

Once

Sim. Descontinuado. Use AvailableNow.

Sim. Descontinuado. Use AvailableNow.

Trigger

Compute clássico

Compute serverless (notebooks e jobs)

RealTime

Sim

Não

ProcessingTime

Sim

Não

AvailableNow

Sim

Sim

Once

Sim. Descontinuado. Use AvailableNow.

Sim. Descontinuado. Use AvailableNow.

Modos de saída

Esta tabela mostra suporte para modos de saída de transmissão estructurada.

Modo de saída

Suportado

update

Sim

append

Sim. O comportamento é idêntico ao de update. A query faz um upsert quando a tabela de destino tem uma chave primária; caso contrário, a query insere. Consulte Comportamento de upsert.

complete

Não

Modo de saída

Suportado

update

Sim

append

Sim. O comportamento é idêntico ao de update. A query faz um upsert quando a tabela de destino tem uma chave primária; caso contrário, a query insere. Consulte Comportamento de upsert.

complete

Não

Limitações

  • Para um banco de dados PostgreSQL externo conectado por meio de uma conexão do Unity Catalog, a tabela de destino já deve existir. O coletor cria automaticamente tabelas ausentes apenas no Lakebase.
  • Lakeflow pipelines are not supported.