Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Use o Structured Streaming para escrever no Lakebase ou em um banco de dados externo PostgreSQL com batching embutido, tentativas automáticas e autenticação gerenciada pelo workspace.
Quando usar o coletor do Lakebase
Use o Lakebase sink para gravações de baixa latência em streaming no Lakebase ou em um banco de dados externo PostgreSQL. Esse coletor não exige que você implemente funções foreach personalizadas para lidar com lotes, gerenciamento de conexões e tratamento de erros.
Os casos de uso comuns incluem:
- Atualize os bancos de dados de aplicativos em tempo real para dashboards operacionais ou recursos voltados para o cliente.
- Sincronizar dados de alteração contínua, como resultados de streaming agregados ou filtrados, em um banco de dados transacional.
- Grave a saída de uma consulta de Structured Streaming em uma tabela do Lakebase com latência inferior a um segundo usando o modo em tempo real.
Para sincronizar dados do Lakebase com tabelas do Delta Lake no Lakehouse, na direção inversa, consulte Lakebase Change Data Feed.
Requisitos
-
Databricks Runtime 18 LTS e superiores.
- Conexões externas com PostgreSQL exigem que você use o Databricks Runtime 19 ou superior e adira à versão prévia do JDBC personalizado no UC Compute.
- Tipos de dados de intervalo exigem o uso do Databricks Runtime 19 ou superior.
- Computação clássica com modos de acesso dedicados ou padrão, ou computação sem servidor para notebooks ou tarefas. Em computação serverless, use
Trigger.AvailableNow(). Consulte Streaming na computação sem servidor. - Um banco de dados Lakebase, ou uma conexão do Unity Catalog para um banco de dados externo PostgreSQL.
Requisitos de identificador
Para todos os alvos, a Databricks recomenda usar nomes de esquemas, tabelas, colunas e colunas de chave primária que começam com uma letra ou sublinhado e contêm apenas letras, números e sublinhados. O sumidouro aplica esses requisitos ao criar automaticamente uma tabela Lakebase. Para usar identificadores que não atendam a esses requisitos, crie a tabela de destino antes de iniciar a consulta.
Conectar a um banco de dados
O coletor lakebase dá suporte aos seguintes métodos de conexão:
Tabelas do Lakebase registradas no Catálogo do Unity
Para tabelas do Lakebase registradas no Catálogo do Unity, o conector gerencia automaticamente as credenciais e usa a identidade do usuário ou da entidade de serviço que executa a consulta. Se a tabela não existir, o conector criará a tabela.
Para registrar um banco de dados do Lakebase com o Catálogo do Unity, consulte Registrar um banco de dados lakebase no Catálogo do Unity.
Para escrever em uma tabela Lakebase, use o .toTable() método 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>")
)
Scala
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 marcadores:
-
<catalog>.<schema>.<table>: o nome totalmente qualificado da tabela de destino. O catálogo do Unity Catalogcatalogé o catálogo 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 criará. -
<primary-key-columns>: opcional. Uma lista separada por vírgulas de todas as colunas na chave primária da tabela de destino, por exemploidouuser_id,event_type. Veja 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 objetos de nuvem. O local deve ser o armazenamento para o qual você pode gravar, não o disco local e deve ser exclusivo para cada consulta de streaming. Isso é independente da tabela de destino. Consulte pontos de verificação do fluxo estruturado.
Para configurações opcionais, como batchsize e batchinterval, veja as opções do sink do PostgreSQL.
Tabelas do Lakebase não registradas no Catálogo do Unity
Para tabelas do Lakebase não registradas no Catálogo do Unity, o conector gerencia automaticamente as credenciais e usa a identidade do usuário ou da entidade de serviço que executa a consulta. Se a tabela não existir, o conector criará 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()
)
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()
Substitua os seguintes marcadores:
-
<project-id>.<branch-id>.<endpoint-id>: seu endpoint do Lakebase. Localize os três valores no nome do recurso no menu Obter ID da guia Computação , que tem o formatoprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Consulte identificadores de computação. -
<database>: opcional. O nome do banco de dados PostgreSQL alvo. Usadatabricks_postgrescomo padrão. Consulte Gerenciar bancos de dados. -
<schema>.<table>: a tabela de destino no formatoschema.table. Se você omitir o esquema, o coletor usará o esquemapublic. Para criação automática de tabelas, use identificadores que começam com uma letra ou sublinhado e contêm apenas letras, números e sublinhados. -
<primary-key-columns>: opcional. Uma lista separada por vírgulas de todas as colunas na chave primária da tabela de destino, por exemploidouuser_id,event_type. Veja 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 objetos de nuvem. O local deve ser o armazenamento para o qual você pode gravar, não o disco local e deve ser exclusivo para cada consulta de streaming. Isso é independente da tabela de destino. Consulte pontos de verificação do fluxo estruturado.
Para configurações opcionais, como batchsize e batchinterval, veja as opções do sink do PostgreSQL.
PostgreSQL externo com credenciais do Catálogo Unity
Importante
Esse recurso está em Visualização Pública. Administradores do espaço de trabalho podem controlar o acesso ao JDBC personalizado no UC Compute na página Pré-visualizações. Consulte Gerenciar visualizações do Azure Databricks.
Use uma conexão com o Unity Catalog para autenticar em um banco de dados externo PostgreSQL sem armazenar credenciais no seu código. A tabela alvo já deve existir.
Crie uma conexão do tipo POSTGRESQL, veja Criar uma conexão. O usuário ou a entidade de serviço que executa a consulta deve ter USE CONNECTION na conexão.
Para escrever na tabela do PostgreSQL, use as opções databricks.connection, database e 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()
Substitua os seguintes marcadores:
-
<connection-name>: O nome da conexão do Unity Catalog. -
<database>: O nome do banco de dados PostgreSQL alvo. -
<schema>.<table>: A tabela de destino existente no formatoschema.table. Se você omitir o esquema, o coletor usará o esquemapublic. -
<primary-key-columns>: opcional. Uma lista separada por vírgulas de todas as colunas na chave primária da tabela de destino, por exemploidouuser_id,event_type. Veja 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 objetos de nuvem. O local deve ser o armazenamento para o qual você pode gravar, não o disco local e deve ser exclusivo para cada consulta de streaming. Isso é independente da tabela de destino. Consulte pontos de verificação do fluxo estruturado.
Conexões PostgreSQL sempre usam TLS. A verificação de certificados segue as configurações na conexão do Unity Catalog, que você escolhe ao criar a conexão:
-
Certificado do servidor de confiança: Quando selecionado, a conexão usa
sslmode=require, que criptografa a conexão sem verificar o certificado do servidor. -
Certificado de servidor fornecido pelo usuário: Fornecer um certificado de servidor codificado em PEM para usar
sslmode=verify-fullquando o certificado do servidor de confiança não for selecionado. Se você não fornecer um certificado, a conexão usasslmode=verify-fullcom o repositório de confiança padrão da JVM.
Opções de configuração
O coletor retorna um erro ao encontrar opções não reconhecidas, JDBC_STREAMING_SINK_INVALID_OPTIONS.
Para as opções de configuração do sink, incluindo as opções comuns e as opções para cada método de conexão, veja as opções do sink do PostgreSQL.
Mapeamentos de tipo de dados
O coletor verifica se cada coluna do DataFrame é compatível com a coluna de destino correspondente antes de gravar em uma tabela existente do Lakebase ou em uma tabela externa do PostgreSQL.
A tabela a seguir contém tipos suportados no Databricks Runtime 18 LTS e superiores:
| Tipo de Spark | Tipo de tabela Lakebase criado automaticamente | Tipos compatíveis em tabelas 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 |
A tabela a seguir contém os tipos suportados no Databricks Runtime 19 e superior:
| Tipo de Spark | Tipo de tabela Lakebase criado automaticamente | Tipos compatíveis em tabelas PostgreSQL existentes |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Comportamento de upsert
A upsertkey opção identifica as colunas principais da tabela de destino. Para uma tabela existente, as colunas em upsertkey devem corresponder exatamente à chave primária da tabela. Se você omitir a opção, o sink lê a chave primária da tabela. Para uma tabela Lakebase que o sumidouro cria, upsertkey define a chave primária. Se você omitir essa opção, o sink cria a tabela sem uma chave primária.
Quando a tabela de destino tem uma chave primária, o sink faz upsert usando a sintaxe INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... do PostgreSQL. Quando a tabela de destino não possui chave primária, o sink executa inserções. O modo de saída de uma consulta não tem efeito nesse comportamento.
Todas as colunas de chave primária devem estar presentes no DataFrame e usar tipos comparáveis, como tipos numéricos ou de texto.
Ajuste de desempenho
Lotes e contrapressão
Uma descarga é acionada quando uma das condições é atendida:
- O buffer atinge
batchsizelinhas, cujo padrão é1000. - A idade do buffer excede
batchinterval, cujo padrão é100 milliseconds.
Quando o banco de dados não consegue acompanhar a taxa de dados de entrada, o coletor propaga contrapressão upstream para a origem.
Diretrizes de latência e taxa de transferência:
- Para cargas de trabalho de baixa latência com modo em tempo real, reduza
batchintervalpara garantir um tempo máximo menor antes da liberação. Veja Conceitos em modo em tempo real para conceitos e exemplos de modo em tempo real para um exemplo de código. - Para cargas de trabalho de alta taxa de transferência, aumente
batchsizepara reduzir a sobrecarga de cada transação.
Comportamento de conexão
O coletor usa pooling de conexões nos executores. Por padrão, cada tarefa usa uma conexão de banco de dados.
O Databricks recomenda que você use o valor padrão da tarefa 1 para cada conexão. Se você 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 baixo de conexões, reduza o número de partições de shuffle ou aumente spark.databricks.sql.streaming.jdbc.tasksPerConnection.
O coletor repete automaticamente erros JDBC transitórios, 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 com suporte
Gatilhos
Esta tabela mostra suporte para tipos de gatilhos de Streaming Estruturado em computação clássica e serverless:
| Gatilho | Computação clássica | Computação sem servidor (notebooks e tarefas) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Sim. Preterido. Use AvailableNow. |
Sim. Preterido. Use AvailableNow. |
Modos de saída
Esta tabela mostra o suporte para modos de saída de Streaming Estruturado:
| Modo de saída | Supported |
|---|---|
update |
Yes |
append |
Sim. O comportamento é idêntico a update. A consulta executa upsert quando a tabela de destino tem uma chave primária; caso contrário, a consulta insere. Veja comportamento de Upsert. |
complete |
No |
Limitações
- Para um banco de dados PostgreSQL externo conectado por meio de uma conexão com o Unity Catalog, a tabela alvo já deve existir. O sink cria automaticamente tabelas ausentes apenas no Lakebase.
- Não há suporte a pipelines do Lakeflow.