create_table

Important

Ez a funkció bétaverzióban érhető el.

A folyamat függvényével create_table() létrehozhat egy felügyelt táblát, amelyet egy vagy több append_flow deklaráció ír. Párosítsa a create_table() hívást egy vagy több @append_flow(target=...) dekoratőrrel, amelyek a táblába írnak. Több folyamat is megcélzhatja ugyanazt a felügyelt táblát.

Az SQL-ekvivalensről lásd:CREATE TABLE ... FLOW.

Syntax

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

Paraméter Típus Description
name str Required. A tábla neve.
comment str A táblázat leírása.
spark_conf dict A lekérdezés végrehajtásához szükséges Spark-konfigurációk listája.
table_properties dict A dicttáblázati tulajdonságok halmaza a táblázathoz.
partition_cols list A tábla particionálásához használandó egy vagy több oszlop listája.
path str A táblaadatok tárolási helye. Ha nincs beállítva, használja a táblát tartalmazó séma felügyelt tárolási helyét.
schema str vagy StructType A tábla sémadefiníciója. A sémák definiálhatók SQL DDL karakterláncként vagy Python StructType-ként.
\, \, \ dict A tábla adatminőségi korlátozásai. Ugyanazt a viselkedést biztosítja, és ugyanazt a szintaxist használja, mint az elvárás dekorátor függvények, de paraméterként implementálva. Lásd az elvárásokat.
cluster_by list Engedélyezze a "liquid clustering" funkciót a táblában, és határozza meg a fürtözési kulcsként használni kívánt oszlopokat. Lásd: Táblákhoz folyékony klaszterezés használata.
cluster_by_auto bool Automatikus folyadékklaszterezés engedélyezése a táblázatban. A kezdeti fürtözési kulcsok definiálásához kombinálható cluster_by . Lásd: Automatikus folyadékfürtözés.
row_filter str (Nyilvános előzetes verzió) A záradék a sorok szűrésére a táblában. Lásd: Táblázatok közzététele sorszűrőkkel és oszlopmaszkokkal.
private bool Amikor Truea rendszer létrehoz egy privát táblát, amely nincs közzétéve a katalógusban, és csak a folyamaton belül érhető el. Alapértelmezett érték: False.

Limitations

  • A felügyelt táblák nem támogatják a változási adatrögzítés (CDC) változásfolyamatait. create_auto_cdc_flow() vagy create_auto_cdc_from_snapshot_flow() egy felügyelt tábla megcélzása meghiúsul. CDC-célokhoz használja a create_streaming_table() elemet.
  • A felügyelt táblák csak append_flowazokat támogatják. A cserefolyamatok (replace_flow / FLOW ... REPLACE WHERE) nem támogatottak.
  • A felügyelt táblák csak a Unity Catalogtal rendelkező folyamatokban támogatottak.
  • Egy felügyelt tábla meglévő streamelési táblájának nevét nem használhatja újra.

Példa

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