Aller au contenu principal

Table

Le décorateur @table peut être utilisé pour définir des tables de streaming dans un pipeline.

Pour définir une table de streaming, appliquez @table à une query qui effectue une lecture en streaming sur une source de données ou utilisez la fonction create_streaming_table().

remarque

Dans l’ancien module dlt, l’opérateur @table était utilisé pour créer à la fois des tables de streaming et des vues matérialisées. L'opérateur @table dans le module pyspark.pipelines fonctionne toujours de cette manière, mais Databricks recommande d'utiliser l'opérateur @materialized_view pour créer des vues matérialisées.

Syntaxe

Python
from pyspark import pipelines as dp

@dp.table(
name="<name>",
comment="<comment>",
spark_conf={&quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;, &quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;},
table_properties={&quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;, &quot;&lt;key&gt;&quot; : &quot;&lt;value&gt;&quot;},
path="<storage-location-path>",
partition_cols=["<partition-column>", "<partition-column>"],
cluster_by_auto = False,
cluster_by = ["<clustering-column>", "<clustering-column>"],
schema="schema-definition",
row_filter = "row-filter-clause",
private = False)
@dp.expect(...)
def <function-name>():
return (<query>)

parameter

@dp.expect() ** ** est une clause d'attente facultative des LakeFlow Pipelines. Vous pouvez inclure plusieurs attentes. Consultez Attentes.

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

name

str

Le nom de la table. Si non fourni, le nom de la fonction default.

comment

str

Une description pour la table.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

table_properties

dict

Un dict de propriétés de table pour la table.

path

str

Un emplacement de stockage pour les données de table. Si non défini, utilisez l'emplacement de stockage géré pour le schéma contenant la table.

partition_cols

list

Une liste d'une ou plusieurs colonnes à utiliser pour le partitionnement de la table.

cluster_by_auto

bool

Activez le clustering liquide automatique sur la table. Cela peut être combiné avec cluster_by et définir les colonnes à utiliser comme clés de clustering initiales, suivi d'un monitoring et de mises à jour automatiques de la sélection des clés basées sur la charge de travail. Voir clustering liquide automatique.

cluster_by

list

Activez le clustering liquide sur la table et définissez les colonnes à utiliser comme clés de clustering. Voir Utiliser le clustering liquide pour les tables.

schema

str OU StructType

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 StructType Python. Pour les propriétés de colonne prises en charge dans la chaîne DDL, consultez la section Parameters de CREATE STREAMING TABLE.

private

bool

Créez une table, mais ne publiez pas la table dans le metastore. Cette table est disponible pour le pipeline mais n'est pas accessible en dehors du pipeline. Les tables privées persistent pendant toute la durée de vie du pipeline.

La default est False.

Les tables privées ont été précédemment créées avec le paramètre temporary.

row_filter

str

(Aperçu public) Une clause de filtre de ligne pour la table. Consultez Publier des tables avec des filtres de ligne et des masques de colonne.

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

name

str

Le nom de la table. Si non fourni, le nom de la fonction default.

comment

str

Une description pour la table.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

table_properties

dict

Un dict de propriétés de table pour la table.

path

str

Un emplacement de stockage pour les données de table. Si non défini, utilisez l'emplacement de stockage géré pour le schéma contenant la table.

partition_cols

list

Une liste d'une ou plusieurs colonnes à utiliser pour le partitionnement de la table.

cluster_by_auto

bool

Activez le clustering liquide automatique sur la table. Cela peut être combiné avec cluster_by et définir les colonnes à utiliser comme clés de clustering initiales, suivi d'un monitoring et de mises à jour automatiques de la sélection des clés basées sur la charge de travail. Voir clustering liquide automatique.

cluster_by

list

Activez le clustering liquide sur la table et définissez les colonnes à utiliser comme clés de clustering. Voir Utiliser le clustering liquide pour les tables.

schema

str OU StructType

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 StructType Python. Pour les propriétés de colonne prises en charge dans la chaîne DDL, consultez la section Parameters de CREATE STREAMING TABLE.

private

bool

Créez une table, mais ne publiez pas la table dans le metastore. Cette table est disponible pour le pipeline mais n'est pas accessible en dehors du pipeline. Les tables privées persistent pendant toute la durée de vie du pipeline.

La default est False.

Les tables privées ont été précédemment créées avec le paramètre temporary.

row_filter

str

(Aperçu public) Une clause de filtre de ligne pour la table. Consultez Publier des tables avec des filtres de ligne et des masques de colonne.

La spécification d'un schéma est facultative et peut être effectuée avec PySpark StructType ou DDL SQL. Lorsque vous spécifiez un schéma, vous pouvez inclure en option des colonnes générées, des masques de colonne, ainsi que des clés primaires et étrangères. Voir :

Exemples

Python
from pyspark import pipelines as dp

# Specify a schema
sales_schema = StructType([
StructField("customer_id", StringType(), True),
StructField("customer_name", StringType(), True),
StructField("number_of_line_items", StringType(), True),
StructField("order_datetime", StringType(), True),
StructField("order_number", LongType(), True)]
)
@dp.table(
comment="Raw data on sales",
schema=sales_schema)
def sales():
return ("...")

# Specify a schema with SQL DDL, use a generated column, and set clustering columns
@dp.table(
comment="Raw data on sales",
schema="""
customer_id STRING,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime))
""",
cluster_by = ["order_day_of_week", "customer_id"])
def sales():
return ("...")

# Specify a schema with an identity column
@dp.table(
comment="Raw data on sales",
schema="""
order_id BIGINT GENERATED ALWAYS AS IDENTITY,
customer_name STRING,
order_datetime STRING
""")
def sales():
return ("...")

# Use automatic liquid clustering to let Databricks choose the clustering columns
@dp.table(
comment="Raw data on sales",
cluster_by_auto=True)
def sales():
return ("...")

# Specify partition columns
@dp.table(
comment="Raw data on sales",
schema="""
customer_id STRING,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime))
""",
partition_cols = ["order_day_of_week"])
def sales():
return ("...")

# Specify table constraints
@dp.table(
schema="""
customer_id STRING NOT NULL PRIMARY KEY,
customer_name STRING,
number_of_line_items STRING,
order_datetime STRING,
order_number LONG,
order_day_of_week STRING GENERATED ALWAYS AS (dayofweek(order_datetime)),
CONSTRAINT fk_customer_id FOREIGN KEY (customer_id) REFERENCES main.default.customers(customer_id)
""")
def sales():
return ("...")

# Specify a row filter and column mask
@dp.table(
schema="""
id int COMMENT 'This is the customer ID',
name string COMMENT 'This is the customer full name',
region string,
ssn string MASK catalog.schema.ssn_mask_fn USING COLUMNS (region)
""",
row_filter = "ROW FILTER catalog.schema.us_filter_fn ON (region, name)")
def sales():
return ("...")