Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Cette fonctionnalité est en version bêta.
Utilisez la create_table() fonction dans un pipeline pour créer une table managée, écrite par une ou plusieurs déclarations de append_flow . Associez l’appel create_table() à un ou plusieurs @append_flow(target=...) décorateurs qui écrivent dans la table. Plusieurs flux peuvent cibler la même table managée.
Pour l’équivalent SQL, consultez 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ètre | Catégorie | Description |
|---|---|---|
name |
str |
Required. Nom de la table. |
comment |
str |
Description de la table. |
spark_conf |
dict |
Liste des configurations Spark pour l’exécution de cette requête. |
table_properties |
dict |
Une dict de propriétés de table pour la table. |
partition_cols |
list |
Liste d’une ou de plusieurs colonnes à utiliser pour partitionner la table. |
path |
str |
Emplacement de stockage pour les données de table. Si ce n’est pas le cas, utilisez l’emplacement de stockage managé pour le schéma contenant la table. |
schema |
str ou StructType |
Définition de schéma pour la table. Les schémas peuvent être définis en tant que chaîne SQL DDL ou avec StructTypePython. |
expect_all, , expect_all_or_dropexpect_all_or_fail |
dict |
Contraintes de qualité des données pour la table. Fournit le même comportement et utilise la même syntaxe que les fonctions de décorateur attendues, mais implémentées en tant que paramètre. Voir les attentes. |
cluster_by |
list |
Activez le clustering liquide sur la table et définissez les colonnes à utiliser comme clés de clustering. Consultez Utilisation de Liquid Clustering pour les tables. |
cluster_by_auto |
bool |
Activez le clustering liquide automatique sur la table. Peut être combiné avec cluster_by pour définir les clés de clustering initiales. Consultez le regroupement automatique de liquide. |
row_filter |
str |
(Préversion publique) Clause de filtre de ligne pour la table. Consultez Publier des tables avec des filtres de lignes et des masques de colonne. |
private |
bool |
Lorsque True, crée une table privée qui n’est pas publiée dans le catalogue et est accessible uniquement dans le pipeline. La valeur par défaut est False. |
Limitations
- Les tables managées ne prennent pas en charge les flux de modification de capture de données modifiées (CDC).
create_auto_cdc_flow()oucreate_auto_cdc_from_snapshot_flow()le ciblage d’une table managée échoue. Utilisez create_streaming_table() pour les cibles CDC. - Les tables managées prennent uniquement
append_flowen charge . Les flux de remplacement (replace_flow/FLOW ... REPLACE WHERE) ne sont pas pris en charge. - Les tables managées sont prises en charge uniquement dans les pipelines avec le catalogue Unity.
- Vous ne pouvez pas réutiliser le nom d’une table de diffusion en continu existante pour une table managé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")