create_table

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() ou create_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")