create_table
Bêta
Cette fonctionnalité est en Bêta.
Utilisez la fonction create_table() dans un pipeline pour créer une table gérée, écrite par une ou plusieurs déclarations append_flow. Associez l'appel create_table() avec un ou plusieurs décorateurs @append_flow(target=...) qui écrivent dans la table. Plusieurs flux peuvent cibler la même table gérée.
Pour l'équivalent SQL, consultez CREATE TABLE... FLOW.
Syntaxe
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
)
parameter
parameter | Type | Description |
|---|---|---|
|
| Obligatoire. Le nom de la table. |
|
| Une description pour la table. |
|
| Une liste de configurations Spark pour l'exécution de cette query. |
|
| Une |
|
| Une liste d'une ou plusieurs colonnes à utiliser pour partitionner la table. |
|
| Un emplacement de stockage pour les données de table. Si aucune valeur n'est spécifiée, utilisez l'emplacement de stockage géré pour le schéma contenant la table. |
|
| Une définition de schéma pour la table. Les schémas peuvent être définis comme une chaîne DDL SQL ou avec un Python |
|
| Contraintes de qualité des données pour la table. Offre le même comportement et utilise la même syntaxe que les fonctions de décorateur d'attente, mais implémenté en tant que paramètre. Consultez Attentes. |
|
| Activez le liquid clustering sur la table et définissez les colonnes à utiliser comme clés de clustering. Voir Utiliser le liquid clustering pour les tables. |
|
| Activez le clustering liquide automatique sur la table. Peut être combiné avec |
|
| (Aperçu public) Une clause de filtre de ligne pour la table. Voir Publier des tables avec des filtres de ligne et des masques de colonne. |
|
| Lorsque |
Limitations
- Les tables gérées ne prennent pas en charge les flux de modifications de la capture des données modifiées (CDC).
create_auto_cdc_flow()oucreate_auto_cdc_from_snapshot_flow()ciblant une table gérée échoue. Utilisez create_streaming_table() pour les cibles CDC. - Les tables gérées ne prennent en charge que
append_flow. Les flux de remplacement (replace_flow/FLOW ... REPLACE WHERE) ne sont pas pris en charge. - Les tables gérées ne sont prises en charge que dans les pipelines avec Unity Catalog.
- Vous ne pouvez pas réutiliser le nom d'une table de streaming existante pour une table gérée.
Exemple
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")