create_table

Important

Dieses Feature befindet sich in der Betaversion.

Verwenden Sie die create_table() Funktion in einer Pipeline, um eine verwaltete Tabelle zu erstellen, die von einer oder mehreren append_flow Deklarationen geschrieben wurde. Koppeln Sie den create_table() Anruf mit einem oder @append_flow(target=...) mehreren Dekoratoren, die in die Tabelle schreiben. Mehrere Flüsse können auf dieselbe verwaltete Tabelle abzielen.

Informationen zur SQL-Entsprechung finden Sie unter 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

Parameter Typ Description
name str Required. Der Tabellenname.
comment str Eine Beschreibung für die Tabelle.
spark_conf dict Eine Liste der Spark-Konfigurationen für die Ausführung dieser Abfrage.
table_properties dict Eine dict von Tabelleneigenschaften für die Tabelle
partition_cols list Eine Liste mit einer oder mehreren Spalten, die für die Partitionierung der Tabelle verwendet werden sollen.
path str Ein Speicherort für Tabellendaten. Wenn sie nicht festgelegt ist, verwenden Sie den verwalteten Speicherort für das Schema, das die Tabelle enthält.
schema str oder StructType Eine Schemadefinition für die Tabelle. Schemas können als SQL-DDL-Zeichenfolge oder mit Python StructType definiert werden
expect_all, expect_all_or_dropexpect_all_or_fail dict Datenqualitätseinschränkungen für die Tabelle. Bietet dasselbe Verhalten und verwendet dieselbe Syntax wie Erwartungsdekoratorfunktionen, ist jedoch als Parameter implementiert Siehe Erwartungen.
cluster_by list Aktivieren des Liquid Clustering für die Tabelle und Definieren der Spalten, die als Clusterschlüssel verwendet werden sollen. Siehe Verwenden von Flüssigclustering für Tabellen.
cluster_by_auto bool Aktivieren Sie die automatische Flüssigkeitsgruppierung auf dem Tisch. Kann kombiniert werden, cluster_by um die anfänglichen Clusteringschlüssel zu definieren. Siehe Automatische Flüssigkeitsclusterung.
row_filter str (Öffentliche Vorschau) Eine Zeilenfilterklausel für die Tabelle. Siehe Veröffentlichen von Tabellen mit Zeilenfiltern und Spaltenmasken.
private bool Wenn True, erstellt eine private Tabelle, die nicht im Katalog veröffentlicht wird und nur innerhalb der Pipeline zugänglich ist. Wird standardmäßig auf False festgelegt.

Einschränkungen

  • Verwaltete Tabellen unterstützen änderungsdatenerfassungs-Änderungsflüsse (CDC) nicht. create_auto_cdc_flow() oder create_auto_cdc_from_snapshot_flow() die Ausrichtung auf eine verwaltete Tabelle schlägt fehl. Verwenden Sie create_streaming_table() für CDC-Ziele.
  • Verwaltete Tabellen unterstützen nur append_flow. Ersetzungsflüsse (replace_flow / FLOW ... REPLACE WHERE) werden nicht unterstützt.
  • Verwaltete Tabellen werden nur in Pipelines mit Unity-Katalog unterstützt.
  • Sie können den Namen einer vorhandenen Streamingtabelle für eine verwaltete Tabelle nicht wiederverwenden.

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