Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Use o Structured Streaming para escrever no Lakebase ou numa base de dados PostgreSQL externa com batching incorporado, tentativas automáticas e autenticação gerida pelo espaço de trabalho.
Quando usar o sumidouro Lakebase
Use o sumidouro do Lakebase para escritas de streaming de baixa latência para o Lakebase ou para uma base de dados PostgreSQL externa. Este sink não exige que implementes funções personalizadas foreach para gerir batching, gestão de ligações e gestão de erros.
Os casos de uso comuns incluem:
- Atualize as bases de dados das aplicações em tempo real para dashboards operacionais ou funcionalidades voltadas para o cliente.
- Sincronizar dados em constante mudança, como resultados agregados ou filtrados de streaming, numa base de dados transacional.
- Escreva a saída de uma consulta de Streaming Estruturado numa tabela Lakebase com latência inferior a um segundo usando o modo em tempo real.
Para sincronizar dados do Lakebase para as tabelas Delta Lake no Lakehouse, na direção inversa, veja Lakebase Change Data Feed.
Requisitos
-
Databricks Runtime 18 LTS e superiores.
- As ligações PostgreSQL externas exigem que utilize o Databricks Runtime 19 e superiores e opte pelo JDBC personalizado na pré-visualização do UC Compute .
- Os tipos de dados de intervalo requerem a utilização de 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. Na computação sem servidor, use
Trigger.AvailableNow(). Veja Transmissão na computação sem servidor. - Uma base de dados Lakebase, ou uma ligação Unity Catalog a uma base de dados PostgreSQL externa.
Requisitos de identificador
Para todos os alvos, a Databricks recomenda a utilização de nomes de esquemas, tabelas, colunas e colunas de chave primária que comecem por uma letra ou carácter de sublinhado e que contenham apenas letras, números e caracteres de sublinhado. O sumidouro faz cumprir estes requisitos quando cria automaticamente uma tabela Lakebase. Para usar identificadores que não cumpram estes requisitos, crie a tabela de destino antes de iniciar a consulta.
Ligar a uma base de dados
O sumidouro Lakebase suporta os seguintes métodos de ligação:
Tabelas Lakebase registadas no Unity Catalog
Para tabelas Lakebase registadas no Unity Catalog, o conector gere automaticamente as credenciais e utiliza a identidade do utilizador ou do principal do serviço que executa a consulta. Se a tabela não existir, o conector cria a tabela.
Para registar uma base de dados Lakebase no Unity Catalog, consulte Registar uma base de dados Lakebase no Unity Catalog.
Para escrever numa 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 placeholders:
-
<catalog>.<schema>.<table>: O nome totalmente qualificado da tabela alvo. Estecatalogé o catálogo do Catálogo Unity que criou quando registou a base de dados Lakebase, veja Registar uma base de dados Lakebase no Catálogo Unity. Se a tabela não existir, o conector cria-a. -
<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. Ver comportamento de Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho para um volume do Unity Catalog onde a consulta armazena o ponto de verificação. Também pode usar um URI de armazenamento de objetos na cloud. A localização deve ser armazenamento onde possas escrever, não disco local, e deve ser única para cada consulta de streaming. Isto é independente da tabela alvo. Consulte Pontos de verificação de streaming estruturado.
Para configurações opcionais, como batchsize e batchinterval, veja as opções do sink do PostgreSQL.
Tabelas Lakebase não registadas no Unity Catalog
Para tabelas Lakebase não registadas no Unity Catalog, o conector gere automaticamente as credenciais e utiliza a identidade do utilizador ou principal do serviço que executa a consulta. Se a tabela não existir, o conector cria a tabela.
Para escrever para uma tabela 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 placeholders:
-
<project-id>.<branch-id>.<endpoint-id>: O seu endpoint Lakebase. Encontre os três valores no nome do Recurso no menu Obter ID do separador Computes , que tem o formatoprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Ver Identificadores de computação. -
<database>: Opcional. O nome da base de dados PostgreSQL alvo. O valor padrão édatabricks_postgres. Consulte Gerir bases de dados. -
<schema>.<table>: A tabela alvo emschema.tableformato. Se omitir o esquema, o sink utiliza o esquemapublic. Para a criação automática de tabelas, utilize identificadores que começam com uma letra ou sublinhado e contenham 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. Ver comportamento de Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho para um volume do Unity Catalog onde a consulta armazena o ponto de verificação. Também pode usar um URI de armazenamento de objetos na cloud. A localização deve ser armazenamento onde possas escrever, não disco local, e deve ser única para cada consulta de streaming. Isto é independente da tabela alvo. Consulte Pontos de verificação de streaming estruturado.
Para configurações opcionais, como batchsize e batchinterval, veja as opções do sink do PostgreSQL.
PostgreSQL externo com credenciais do Unity Catalog
Importante
Este recurso está no Public Preview. Os administradores do espaço de trabalho podem controlar o acesso ao JDBC Personalizado no UC Compute a partir da página de Pré-visualizações . Ver Gerir as pré-visualizações de Azure Databricks.
Use uma ligação ao Unity Catalog para autenticar numa base de dados PostgreSQL externa sem armazenar credenciais no seu código. A tabela alvo já deve existir.
Crie uma ligação do tipo POSTGRESQL, veja Criar uma ligação. O utilizador ou principal de serviço que executa a consulta tem de ter USE CONNECTION na ligaçã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 placeholders:
-
<connection-name>: O nome da conexão do Catálogo Unity. -
<database>: O nome da base de dados PostgreSQL alvo. -
<schema>.<table>: A tabela de destino existente no formatoschema.table. Se omitir o esquema, o sink utiliza 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. Ver comportamento de Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Um caminho para um volume do Unity Catalog onde a consulta armazena o ponto de verificação. Também pode usar um URI de armazenamento de objetos na cloud. A localização deve ser armazenamento onde possas escrever, não disco local, e deve ser única para cada consulta de streaming. Isto é independente da tabela alvo. Consulte Pontos de verificação de streaming estruturado.
As ligações PostgreSQL usam sempre TLS. A verificação de certificados segue as definições da ligação do Unity Catalog, que seleciona ao criar a ligação:
-
Certificado de servidor de confiança: Quando selecionado, a ligação usa
sslmode=require, que encripta a ligação sem verificar o certificado do servidor. -
Certificado de servidor fornecido pelo utilizador: Forneça um certificado de servidor codificado em PEM para usar
sslmode=verify-fullquando o certificado do servidor de confiança não for selecionado. Se não fornecer um certificado, a ligação é usadasslmode=verify-fullcom a loja de confiança padrão da JVM.
Opções de configuração
O sumidouro gera um erro para 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 ligação, veja as opções do sink PostgreSQL.
Mapeamentos de tipo de dados
O sink verifica se cada coluna DataFrame é compatível com a sua coluna alvo correspondente antes de escrever numa tabela Lakebase existente ou numa tabela PostgreSQL externa.
A tabela seguinte contém tipos suportados no Databricks Runtime 18 LTS e superiores:
| Tipo de faísca | 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 seguinte contém os tipos suportados no Databricks Runtime 19 e superiores:
| Tipo de faísca | 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-chave primárias da tabela de destino. Para uma tabela existente, as colunas em upsertkey devem corresponder exatamente à chave primária da tabela. Se omitires essa opção, o sink lê a chave primária da tabela. Para uma tabela Lakebase criada pelo sumidouro, upsertkey define a chave primária. Se omitir a opção, o sink cria a tabela sem uma chave primária.
Quando a tabela de destino tem uma chave primária, o destino faz upsert utilizando a sintaxe INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... do PostgreSQL. Quando a tabela alvo não tem chave primária, o sumidouro realiza inserções. O modo de saída de uma consulta não tem efeito neste comportamento.
Todas as colunas-chave primárias devem estar presentes no DataFrame e usar tipos comparáveis, como tipos numéricos ou de cadeia.
Afinação de desempenho
Processamento em lote e retropressão
Uma descarga é acionada quando se verifica uma das condições:
- O buffer atinge
batchsizelinhas, sendo o valor predefinido1000. - A antiguidade do buffer excede
batchinterval, cujo valor por defeito é100 milliseconds.
Quando a base de dados não consegue acompanhar a taxa de dados recebidos, o sumidouro propaga a contrapressão a montante até à fonte.
Orientação sobre latência e débito:
- Para cargas de trabalho de baixa latência com modo de tempo real, diminua
batchintervalpara garantir um tempo máximo mais curto antes da descarga. Consulte conceitos de modo em tempo real para conceitos e exemplos de modo em tempo real para um exemplo de código. - Para cargas de trabalho de alto rendimento, aumente
batchsizepara reduzir a sobrecarga de cada transação.
Comportamento de ligação
O sink utiliza agrupamento de ligações nos executores. Por defeito, cada tarefa utiliza uma ligação à base de dados.
A Databricks recomenda que utilize o valor predefinido de 1 para cada ligação. Se aumentares o número de tarefas para cada ligação, podes causar conflitos na ligação e aumentar as latências para ligações de alto débito.
Para configurar a proporção de tarefas para ligações, defina a spark.databricks.sql.streaming.jdbc.tasksPerConnection configuração do Spark. Se a base de dados alvo tiver um limite baixo de ligações, reduza o número de partições de mistura ou aumente spark.databricks.sql.streaming.jdbc.tasksPerConnection.
O sink tenta automaticamente erros JDBC transitórios, incluindo falhas de ligação, deadlocks e limitação de taxa. Se o sink esgotar o número máximo de tentativas, a consulta falha.
Acionadores suportados e modos de saída
Triggers
Esta tabela mostra suporte para tipos de gatilhos de Streaming Estruturado em computação clássica e serverless:
| Trigger | Computação clássica | Computação serverless (blocos de notas e tarefas) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Sim. Deprecated. Utilize AvailableNow. |
Sim. Deprecated. Utilize 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 faz um upsert quando a tabela de destino tem uma chave primária; caso contrário, a consulta insere. Ver comportamento de Upsert. |
complete |
No |
Limitações
- Para uma base de dados PostgreSQL externa ligada através de uma ligação ao Unity Catalog, a tabela alvo deve já existir. O sumidouro cria automaticamente tabelas em falta apenas no Lakebase.
- Os oleodutos de fluxo de lago não são suportados.