Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
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()odercreate_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")