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.
Importante
Esse recurso está em Beta.
Use a create_table() função em um pipeline para criar uma tabela gerenciada, escrita por uma ou mais declarações append_flow . Emparelhe a create_table() chamada com um ou mais @append_flow(target=...) decoradores que gravam na tabela. Vários fluxos podem ter como destino a mesma tabela gerenciada.
Para o equivalente do SQL, consulte CREATE TABLE ... FLUXO.
Sintaxe
from pyspark import pipelines as dp
dp.create_table(
name = "<table-name>",
comment = "<comment>",
spark_conf={"<key>" : "<value>", "<key>" : "<value>"},
table_properties={"<key>" : "<value>", "<key>" : "<value>"},
partition_cols=["<partition-column>", "<partition-column>"],
path="<storage-location-path>",
schema="schema-definition",
expect_all = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_drop = {"<key>" : "<value>", "<key>" : "<value>"},
expect_all_or_fail = {"<key>" : "<value>", "<key>" : "<value>"},
cluster_by = ["<clustering-column>", "<clustering-column>"],
cluster_by_auto = False,
row_filter = "row-filter-clause",
private = False
)
Parameters
| Parâmetro | Tipo | Descrição |
|---|---|---|
name |
str |
Required. O nome da tabela. |
comment |
str |
Uma descrição da tabela. |
spark_conf |
dict |
Uma lista de configurações do Spark para a execução dessa consulta. |
table_properties |
dict |
Um dict das propriedades da tabela para a tabela. |
partition_cols |
list |
Uma lista de uma ou mais colunas a serem usadas para particionar a tabela. |
path |
str |
Um local de armazenamento para dados de tabela. Se não estiver definido, use o local de armazenamento gerenciado para o esquema que contém a tabela. |
schema |
str ou StructType |
Uma definição de esquema para a tabela. Os esquemas podem ser definidos como uma string DDL do SQL ou com um script Python StructType. |
expect_all, , expect_all_or_dropexpect_all_or_fail |
dict |
Restrições de qualidade de dados para a tabela. Fornece o mesmo comportamento e usa a mesma sintaxe que as funções de decorador de expectativa, mas implementada como um parâmetro. Veja as expectativas. |
cluster_by |
list |
Habilite o agrupamento líquido na tabela e defina as colunas a serem usadas como chaves de agrupamento. Consulte Usar clustering líquido para tabelas. |
cluster_by_auto |
bool |
Habilite o agrupamento automático de líquidos na tabela. Pode ser combinado para cluster_by definir as chaves de clustering iniciais. Consulte clusterização automática de líquidos. |
row_filter |
str |
(Versão prévia pública) Uma cláusula de filtro de linha para a tabela. Consulte Publicar as tabelas com os filtros de linha e as máscaras de coluna. |
private |
bool |
Quando True, cria uma tabela privada que não é publicada no catálogo e só é acessível dentro do pipeline. Usa False como padrão. |
Limitações
- As tabelas gerenciadas não dão suporte a fluxos de alteração CDC (captura de dados de alteração).
create_auto_cdc_flow()oucreate_auto_cdc_from_snapshot_flow()a segmentação de uma tabela gerenciada falha. Use create_streaming_table() para destinos CDC. - Somente tabelas gerenciadas dão suporte
append_flowa . Não há suporte para fluxos de substituição (replace_flow/FLOW ... REPLACE WHERE). - As tabelas gerenciadas têm suporte apenas em pipelines com o Catálogo do Unity.
- Não é possível reutilizar o nome de uma tabela de streaming existente para uma tabela gerenciada.
Example
from pyspark import pipelines as dp
dp.create_table("combined")
@dp.append_flow(target="combined")
def from_a():
return spark.readStream.table("source_a")
@dp.append_flow(target="combined")
def from_b():
return spark.readStream.table("source_b")