Aller au contenu principal

Classe DataFrame

Une collection distribuée de données regroupées en colonnes nommées.

Un DataFrame est l'équivalent d'une table relationnelle dans Spark SQL et peut être créé à l'aide de diverses fonctions dans SparkSession.

important

Un DataFrame ne doit pas être créé directement à l'aide du constructeur.

Prend en charge Spark Connect

Propriétés

Propriété

Description

sparkSession

Renvoie SparkSession qui a créé ce DataFrame.

rdd

Renvoie le contenu en tant que RDD de lignes (mode Classique uniquement).

na

Renvoie un DataFrameNaFunctions pour gérer les valeurs manquantes.

stat

Renvoie un DataFrameStatFunctions pour les fonctions statistiques.

write

Interface pour enregistrer le contenu du DataFrame non-streaming dans un stockage externe.

writeStream

Interface pour enregistrer le contenu du DataFrame en streaming vers un stockage externe.

schema

Renvoie le schéma de ce DataFrame sous forme de StructType.

dtypes

Renvoie tous les noms de colonne et leurs types de données sous forme de liste.

columns

Récupère les noms de toutes les colonnes du DataFrame sous forme de liste.

storageLevel

Obtenez le niveau de stockage actuel du DataFrame.

isStreaming

Retourne True si ce DataFrame contient une ou plusieurs sources qui renvoient continuellement des données à mesure qu'elles arrivent.

executionInfo

Retourne un objet ExecutionInfo après l'exécution de la query.

plot

Retourne un PySparkPlotAccessor pour les fonctions de traçage.

Propriété

Description

sparkSession

Renvoie SparkSession qui a créé ce DataFrame.

rdd

Renvoie le contenu en tant que RDD de lignes (mode Classique uniquement).

na

Renvoie un DataFrameNaFunctions pour gérer les valeurs manquantes.

stat

Renvoie un DataFrameStatFunctions pour les fonctions statistiques.

write

Interface pour enregistrer le contenu du DataFrame non-streaming dans un stockage externe.

writeStream

Interface pour enregistrer le contenu du DataFrame en streaming vers un stockage externe.

schema

Renvoie le schéma de ce DataFrame sous forme de StructType.

dtypes

Renvoie tous les noms de colonne et leurs types de données sous forme de liste.

columns

Récupère les noms de toutes les colonnes du DataFrame sous forme de liste.

storageLevel

Obtenez le niveau de stockage actuel du DataFrame.

isStreaming

Retourne True si ce DataFrame contient une ou plusieurs sources qui renvoient continuellement des données à mesure qu'elles arrivent.

executionInfo

Retourne un objet ExecutionInfo après l'exécution de la query.

plot

Retourne un PySparkPlotAccessor pour les fonctions de traçage.

Méthodes

Visualisation et inspection des données

Méthode

Description

toJSON(use_unicode)

Convertit un DataFrame en un RDD de chaînes de caractères ou un DataFrame.

printSchema(level)

Affiche le schéma en arborescence.

explain(extended, mode)

Imprime les plans (logiques et physiques) dans la console à des fins de debugging.

show(n, truncate, vertical)

Imprime les n premières lignes du DataFrame dans la console.

collect()

Renvoie tous les enregistrements du DataFrame sous forme de liste de lignes.

toLocalIterator(prefetchPartitions)

Renvoie un itérateur qui contient toutes les lignes de ce DataFrame.

take(num)

Retourne les premières num lignes sous forme de liste de lignes.

tail(num)

Renvoie les num dernières lignes sous forme de liste de lignes.

head(n)

Renvoie les n premières lignes.

first()

Renvoie la première ligne en tant que ligne.

count()

Renvoie le nombre de lignes dans ce DataFrame.

isEmpty()

Vérifie si le DataFrame est vide et renvoie une valeur booléenne.

describe(*cols)

Calcule les statistiques de base pour les colonnes numériques et de chaîne.

summary(*statistics)

Calcule les statistiques spécifiées pour les colonnes numériques et de chaîne.

Méthode

Description

toJSON(use_unicode)

Convertit un DataFrame en un RDD de chaînes de caractères ou un DataFrame.

printSchema(level)

Affiche le schéma en arborescence.

explain(extended, mode)

Imprime les plans (logiques et physiques) dans la console à des fins de debugging.

show(n, truncate, vertical)

Imprime les n premières lignes du DataFrame dans la console.

collect()

Renvoie tous les enregistrements du DataFrame sous forme de liste de lignes.

toLocalIterator(prefetchPartitions)

Renvoie un itérateur qui contient toutes les lignes de ce DataFrame.

take(num)

Retourne les premières num lignes sous forme de liste de lignes.

tail(num)

Renvoie les num dernières lignes sous forme de liste de lignes.

head(n)

Renvoie les n premières lignes.

first()

Renvoie la première ligne en tant que ligne.

count()

Renvoie le nombre de lignes dans ce DataFrame.

isEmpty()

Vérifie si le DataFrame est vide et renvoie une valeur booléenne.

describe(*cols)

Calcule les statistiques de base pour les colonnes numériques et de chaîne.

summary(*statistics)

Calcule les statistiques spécifiées pour les colonnes numériques et de chaîne.

Vues temporaires

Méthode

Description

createTempView(name)

Crée une vue temporaire locale avec ce DataFrame.

createOrReplaceTempView(name)

Crée ou remplace une vue temporaire locale avec ce DataFrame.

createGlobalTempView(name)

Crée une vue temporaire globale avec ce DataFrame.

createOrReplaceGlobalTempView(name)

Crée ou remplace une vue temporaire globale en utilisant le nom donné.

Méthode

Description

createTempView(name)

Crée une vue temporaire locale avec ce DataFrame.

createOrReplaceTempView(name)

Crée ou remplace une vue temporaire locale avec ce DataFrame.

createGlobalTempView(name)

Crée une vue temporaire globale avec ce DataFrame.

createOrReplaceGlobalTempView(name)

Crée ou remplace une vue temporaire globale en utilisant le nom donné.

Sélection et projection

Méthode

Description

select(*cols)

Projete un ensemble d'expressions et renvoie un nouveau DataFrame.

selectExpr(*expr)

Projette un ensemble d'expressions SQL et retourne un nouveau DataFrame.

filter(condition)

Filtre les lignes selon la condition donnée.

where(condition)

Alias pour le filtre.

drop(*cols)

Renvoie un nouveau DataFrame sans les colonnes spécifiées.

toDF(*cols)

Renvoie un nouveau DataFrame avec de nouveaux noms de colonne spécifiés.

withColumn(colName, col)

Renvoie un nouveau DataFrame en ajoutant une colonne ou en remplaçant la colonne existante qui a le même nom.

withColumns(*colsMap)

Renvoie un nouveau DataFrame en ajoutant plusieurs colonnes ou en remplaçant les colonnes existantes qui ont les mêmes noms.

withColumnRenamed(existing, new)

Renvoie un nouveau DataFrame en renommant une colonne existante.

withColumnsRenamed(colsMap)

Renvoie un nouveau DataFrame en renommant plusieurs colonnes.

withMetadata(columnName, metadata)

Retourne un nouveau DataFrame en mettant à jour une colonne existante avec des métadonnées.

metadataColumn(colName)

Sélectionne une colonne de métadonnées basée sur son nom de colonne logique et la retourne en tant que colonne.

colRegex(colName)

Sélectionne une colonne basée sur le nom de colonne spécifié comme expression régulière et la renvoie comme Colonne.

Méthode

Description

select(*cols)

Projete un ensemble d'expressions et renvoie un nouveau DataFrame.

selectExpr(*expr)

Projette un ensemble d'expressions SQL et retourne un nouveau DataFrame.

filter(condition)

Filtre les lignes selon la condition donnée.

where(condition)

Alias pour le filtre.

drop(*cols)

Renvoie un nouveau DataFrame sans les colonnes spécifiées.

toDF(*cols)

Renvoie un nouveau DataFrame avec de nouveaux noms de colonne spécifiés.

withColumn(colName, col)

Renvoie un nouveau DataFrame en ajoutant une colonne ou en remplaçant la colonne existante qui a le même nom.

withColumns(*colsMap)

Renvoie un nouveau DataFrame en ajoutant plusieurs colonnes ou en remplaçant les colonnes existantes qui ont les mêmes noms.

withColumnRenamed(existing, new)

Renvoie un nouveau DataFrame en renommant une colonne existante.

withColumnsRenamed(colsMap)

Renvoie un nouveau DataFrame en renommant plusieurs colonnes.

withMetadata(columnName, metadata)

Retourne un nouveau DataFrame en mettant à jour une colonne existante avec des métadonnées.

metadataColumn(colName)

Sélectionne une colonne de métadonnées basée sur son nom de colonne logique et la retourne en tant que colonne.

colRegex(colName)

Sélectionne une colonne basée sur le nom de colonne spécifié comme expression régulière et la renvoie comme Colonne.

Tri et classement

Méthode

Description

sort(*cols, **kwargs)

Renvoie un nouveau DataFrame trié par la ou les colonnes spécifiées.

orderBy(*cols, **kwargs)

Alias de tri.

sortWithinPartitions(*cols, **kwargs)

Renvoie un nouveau DataFrame, chaque partition étant triée par la ou les colonnes spécifiées.

Méthode

Description

sort(*cols, **kwargs)

Renvoie un nouveau DataFrame trié par la ou les colonnes spécifiées.

orderBy(*cols, **kwargs)

Alias de tri.

sortWithinPartitions(*cols, **kwargs)

Renvoie un nouveau DataFrame, chaque partition étant triée par la ou les colonnes spécifiées.

Agrégation et regroupement

Méthode

Description

groupBy(*cols)

Regroupe le DataFrame par les colonnes spécifiées afin que l'agrégation puisse être effectuée sur celles-ci.

rollup(*cols)

Créez un regroupement multidimensionnel pour le DataFrame actuel en utilisant les colonnes spécifiées.

cube(*cols)

Créez un cube multidimensionnel pour le DataFrame actuel à l'aide des colonnes spécifiées.

groupingSets(groupingSets, *cols)

Créer une agrégation multidimensionnelle pour le DataFrame actuel à l’aide des ensembles de regroupement spécifiés.

agg(*exprs)

Agréger sur l'intégralité du DataFrame sans groupes (raccourci pour df.groupBy().agg()).

observe(observation, *exprs)

Définissez des métriques (nommées) à observer sur le DataFrame.

Méthode

Description

groupBy(*cols)

Regroupe le DataFrame par les colonnes spécifiées afin que l'agrégation puisse être effectuée sur celles-ci.

rollup(*cols)

Créez un regroupement multidimensionnel pour le DataFrame actuel en utilisant les colonnes spécifiées.

cube(*cols)

Créez un cube multidimensionnel pour le DataFrame actuel à l'aide des colonnes spécifiées.

groupingSets(groupingSets, *cols)

Créer une agrégation multidimensionnelle pour le DataFrame actuel à l’aide des ensembles de regroupement spécifiés.

agg(*exprs)

Agréger sur l'intégralité du DataFrame sans groupes (raccourci pour df.groupBy().agg()).

observe(observation, *exprs)

Définissez des métriques (nommées) à observer sur le DataFrame.

Jointures

Méthode

Description

join(other, on, how)

Joint à un autre DataFrame, en utilisant l'expression de jointure donnée.

crossJoin(other)

Renvoie le produit cartésien avec un autre DataFrame.

lateralJoin(other, on, how)

Jointures latérales avec un autre DataFrame, en utilisant l'expression de jointure donnée.

Méthode

Description

join(other, on, how)

Joint à un autre DataFrame, en utilisant l'expression de jointure donnée.

crossJoin(other)

Renvoie le produit cartésien avec un autre DataFrame.

lateralJoin(other, on, how)

Jointures latérales avec un autre DataFrame, en utilisant l'expression de jointure donnée.

Opérations

Méthode

Description

union(other)

Retourne un nouveau DataFrame contenant l'union des lignes dans ce DataFrame et un autre.

unionByName(other, allowMissingColumns)

Renvoie un nouveau DataFrame contenant l'union des lignes de celui-ci et d'un autre DataFrame.

intersect(other)

Renvoie un nouveau DataFrame contenant uniquement les lignes présentes à la fois dans ce DataFrame et dans un autre DataFrame.

intersectAll(other)

Renvoie un nouveau DataFrame contenant des lignes de ce DataFrame et d'un autre DataFrame, tout en préservant les doublons.

subtract(other)

Renvoie un nouveau DataFrame contenant les lignes de ce DataFrame mais pas celles d'un autre DataFrame.

exceptAll(other)

Retourne un nouveau DataFrame contenant des lignes de ce DataFrame mais pas de l'autre DataFrame, tout en préservant les doublons.

Méthode

Description

union(other)

Retourne un nouveau DataFrame contenant l'union des lignes dans ce DataFrame et un autre.

unionByName(other, allowMissingColumns)

Renvoie un nouveau DataFrame contenant l'union des lignes de celui-ci et d'un autre DataFrame.

intersect(other)

Renvoie un nouveau DataFrame contenant uniquement les lignes présentes à la fois dans ce DataFrame et dans un autre DataFrame.

intersectAll(other)

Renvoie un nouveau DataFrame contenant des lignes de ce DataFrame et d'un autre DataFrame, tout en préservant les doublons.

subtract(other)

Renvoie un nouveau DataFrame contenant les lignes de ce DataFrame mais pas celles d'un autre DataFrame.

exceptAll(other)

Retourne un nouveau DataFrame contenant des lignes de ce DataFrame mais pas de l'autre DataFrame, tout en préservant les doublons.

Déduplication

Méthode

Description

distinct()

Renvoie un nouveau DataFrame contenant les lignes distinctes dans ce DataFrame.

dropDuplicates(subset)

Renvoie un nouveau DataFrame avec les lignes en double supprimées, en ne considérant que certaines colonnes, le cas échéant.

dropDuplicatesWithinWatermark(subset)

Renvoie un nouveau DataFrame avec les lignes dupliquées supprimées, en tenant éventuellement compte uniquement de certaines colonnes, dans les limites du filigrane.

Méthode

Description

distinct()

Renvoie un nouveau DataFrame contenant les lignes distinctes dans ce DataFrame.

dropDuplicates(subset)

Renvoie un nouveau DataFrame avec les lignes en double supprimées, en ne considérant que certaines colonnes, le cas échéant.

dropDuplicatesWithinWatermark(subset)

Renvoie un nouveau DataFrame avec les lignes dupliquées supprimées, en tenant éventuellement compte uniquement de certaines colonnes, dans les limites du filigrane.

Échantillonnage et fractionnement

Méthode

Description

sample(withReplacement, fraction, seed)

Renvoie un sous-ensemble échantillonné de ce DataFrame.

sampleBy(col, fractions, seed)

Renvoie un échantillon stratifié sans remplacement basé sur la fraction donnée sur chaque strate.

randomSplit(weights, seed)

Divise aléatoirement ce DataFrame avec les pondérations fournies.

Méthode

Description

sample(withReplacement, fraction, seed)

Renvoie un sous-ensemble échantillonné de ce DataFrame.

sampleBy(col, fractions, seed)

Renvoie un échantillon stratifié sans remplacement basé sur la fraction donnée sur chaque strate.

randomSplit(weights, seed)

Divise aléatoirement ce DataFrame avec les pondérations fournies.

Partitionnement

Méthode

Description

coalesce(numPartitions)

Retourne un nouveau DataFrame qui a exactement numPartitions partitions.

repartition(numPartitions, *cols)

Renvoie un nouveau DataFrame partitionné par les expressions de partitionnement données.

repartitionByRange(numPartitions, *cols)

Renvoie un nouveau DataFrame partitionné par les expressions de partitionnement données.

repartitionById(numPartitions, partitionIdCol)

Renvoie un nouveau DataFrame partitionné par l'expression d'ID de partition donnée.

Méthode

Description

coalesce(numPartitions)

Retourne un nouveau DataFrame qui a exactement numPartitions partitions.

repartition(numPartitions, *cols)

Renvoie un nouveau DataFrame partitionné par les expressions de partitionnement données.

repartitionByRange(numPartitions, *cols)

Renvoie un nouveau DataFrame partitionné par les expressions de partitionnement données.

repartitionById(numPartitions, partitionIdCol)

Renvoie un nouveau DataFrame partitionné par l'expression d'ID de partition donnée.

Remodelage

Méthode

Description

unpivot(ids, values, variableColumnName, valueColumnName)

Dépivoter un DataFrame d'un format large à un format long.

melt(ids, values, variableColumnName, valueColumnName)

Alias pour unpivot.

transpose(indexColumn)

Transpose un DataFrame de sorte que les valeurs de la colonne d'index spécifiée deviennent les nouvelles colonnes.

Méthode

Description

unpivot(ids, values, variableColumnName, valueColumnName)

Dépivoter un DataFrame d'un format large à un format long.

melt(ids, values, variableColumnName, valueColumnName)

Alias pour unpivot.

transpose(indexColumn)

Transpose un DataFrame de sorte que les valeurs de la colonne d'index spécifiée deviennent les nouvelles colonnes.

Gestion des données manquantes

Méthode

Description

dropna(how, thresh, subset)

Renvoie un nouveau DataFrame omettant les lignes avec des valeurs nulles ou NaN.

fillna(value, subset)

Renvoie un nouveau DataFrame dont les valeurs nulles sont remplies avec la nouvelle valeur.

replace(to_replace, value, subset)

Renvoie un nouveau DataFrame remplaçant une valeur par une autre valeur.

Méthode

Description

dropna(how, thresh, subset)

Renvoie un nouveau DataFrame omettant les lignes avec des valeurs nulles ou NaN.

fillna(value, subset)

Renvoie un nouveau DataFrame dont les valeurs nulles sont remplies avec la nouvelle valeur.

replace(to_replace, value, subset)

Renvoie un nouveau DataFrame remplaçant une valeur par une autre valeur.

Fonctions statistiques

Méthode

Description

approxQuantile(col, probabilities, relativeError)

Calcule les quantiles approximatifs des colonnes numériques d'un DataFrame.

corr(col1, col2, method)

Calcule la corrélation de deux colonnes d'un DataFrame sous forme de valeur double.

cov(col1, col2)

Calcule la covariance d'échantillon pour les colonnes données, spécifiées par leurs noms.

crosstab(col1, col2)

Calcule une table de fréquences par paires des colonnes données.

freqItems(cols, support)

Recherche d'éléments fréquents pour les colonnes, éventuellement avec de faux positifs.

Méthode

Description

approxQuantile(col, probabilities, relativeError)

Calcule les quantiles approximatifs des colonnes numériques d'un DataFrame.

corr(col1, col2, method)

Calcule la corrélation de deux colonnes d'un DataFrame sous forme de valeur double.

cov(col1, col2)

Calcule la covariance d'échantillon pour les colonnes données, spécifiées par leurs noms.

crosstab(col1, col2)

Calcule une table de fréquences par paires des colonnes données.

freqItems(cols, support)

Recherche d'éléments fréquents pour les colonnes, éventuellement avec de faux positifs.

Opérations de schéma

Méthode

Description

to(schema)

Renvoie un nouveau DataFrame où chaque ligne est rapprochée pour correspondre au schéma spécifié.

alias(alias)

Renvoie un nouveau DataFrame avec un alias défini.

Méthode

Description

to(schema)

Renvoie un nouveau DataFrame où chaque ligne est rapprochée pour correspondre au schéma spécifié.

alias(alias)

Renvoie un nouveau DataFrame avec un alias défini.

Itération

Méthode

Description

foreach(f)

Applique la fonction f à toutes les lignes de ce DataFrame.

foreachPartition(f)

Applique la fonction f à chaque partition de ce DataFrame.

Méthode

Description

foreach(f)

Applique la fonction f à toutes les lignes de ce DataFrame.

foreachPartition(f)

Applique la fonction f à chaque partition de ce DataFrame.

Mise en cache et persistance

Méthode

Description

cache()

Conserve le DataFrame avec le niveau de stockage default (MEMORY_AND_DISK_DESER).

persist(storageLevel)

Définit le niveau de stockage pour persister le contenu du DataFrame entre les Opérations.

unpersist(blocking)

Marque le DataFrame comme non persistant et supprime tous les blocs le concernant de la mémoire et du disque.

Méthode

Description

cache()

Conserve le DataFrame avec le niveau de stockage default (MEMORY_AND_DISK_DESER).

persist(storageLevel)

Définit le niveau de stockage pour persister le contenu du DataFrame entre les Opérations.

unpersist(blocking)

Marque le DataFrame comme non persistant et supprime tous les blocs le concernant de la mémoire et du disque.

Points de contrôle

Méthode

Description

checkpoint(eager)

Renvoie une version avec point de contrôle de ce DataFrame.

localCheckpoint(eager, storageLevel)

Renvoie une version de ce DataFrame avec point de contrôle local.

Méthode

Description

checkpoint(eager)

Renvoie une version avec point de contrôle de ce DataFrame.

localCheckpoint(eager, storageLevel)

Renvoie une version de ce DataFrame avec point de contrôle local.

Opérations de streaming

Méthode

Description

withWatermark(eventTime, delayThreshold)

Définit un filigrane temporel d'événement pour ce DataFrame.

Méthode

Description

withWatermark(eventTime, delayThreshold)

Définit un filigrane temporel d'événement pour ce DataFrame.

Conseils d'optimisation

Méthode

Description

hint(name, *parameters)

Spécifie une indication sur le DataFrame actuel.

Méthode

Description

hint(name, *parameters)

Spécifie une indication sur le DataFrame actuel.

Limites et décalages

Méthode

Description

limit(num)

Limite le nombre de résultats au nombre spécifié.

offset(num)

Retourne un nouveau DataFrame en ignorant les n premières lignes.

Méthode

Description

limit(num)

Limite le nombre de résultats au nombre spécifié.

offset(num)

Retourne un nouveau DataFrame en ignorant les n premières lignes.

Transformations avancées

Méthode

Description

transform(func, *args, **kwargs)

Renvoie un nouveau DataFrame. Syntaxe concise pour enchaîner les transformations personnalisées.

Méthode

Description

transform(func, *args, **kwargs)

Renvoie un nouveau DataFrame. Syntaxe concise pour enchaîner les transformations personnalisées.

Méthodes de conversion

Méthode

Description

toPandas()

Renvoie le contenu de ce DataFrame sous forme de DataFrame Pandas pandas.DataFrame.

toArrow()

Renvoie le contenu de ce DataFrame en tant que pyarrow.Table PyArrow.

pandas_api(index_col)

Convertit le DataFrame existant en un DataFrame pandas-on-Spark.

mapInPandas(func, schema, barrier, profile)

Mappe un itérateur de batches dans le DataFrame actuel à l'aide d'une fonction Python native.

mapInArrow(func, schema, barrier, profile)

Mappe un itérateur de batches dans le DataFrame actuel à l'aide d'une fonction native Python qui est exécutée sur pyarrow.RecordBatch.

Méthode

Description

toPandas()

Renvoie le contenu de ce DataFrame sous forme de DataFrame Pandas pandas.DataFrame.

toArrow()

Renvoie le contenu de ce DataFrame en tant que pyarrow.Table PyArrow.

pandas_api(index_col)

Convertit le DataFrame existant en un DataFrame pandas-on-Spark.

mapInPandas(func, schema, barrier, profile)

Mappe un itérateur de batches dans le DataFrame actuel à l'aide d'une fonction Python native.

mapInArrow(func, schema, barrier, profile)

Mappe un itérateur de batches dans le DataFrame actuel à l'aide d'une fonction native Python qui est exécutée sur pyarrow.RecordBatch.

Écriture de données

Méthode

Description

writeTo(table)

Créez un constructeur de configuration d'écriture pour les sources v2.

mergeInto(table, condition)

Merge un ensemble de mises à jour, d'insertions et de suppressions basées sur une table source dans une table cible.

Méthode

Description

writeTo(table)

Créez un constructeur de configuration d'écriture pour les sources v2.

mergeInto(table, condition)

Merge un ensemble de mises à jour, d'insertions et de suppressions basées sur une table source dans une table cible.

Comparaison de DataFrame

Méthode

Description

sameSemantics(other)

Renvoie True lorsque les plans de requête logiques dans les deux DataFrames sont égaux.

semanticHash()

Renvoie un code de hachage du plan de requête logique de ce DataFrame.

Méthode

Description

sameSemantics(other)

Renvoie True lorsque les plans de requête logiques dans les deux DataFrames sont égaux.

semanticHash()

Renvoie un code de hachage du plan de requête logique de ce DataFrame.

Métadonnées et informations sur les fichiers

Méthode

Description

inputFiles()

Renvoie un instantané au mieux des fichiers qui composent ce DataFrame.

Méthode

Description

inputFiles()

Renvoie un instantané au mieux des fichiers qui composent ce DataFrame.

Fonctionnalités SQL avancées

Méthode

Description

isLocal()

Renvoie `Vrai` si les méthodes `collect` et `take` peuvent être exécutées localement.

asTable()

Convertit le DataFrame en un objet TableArg, qui peut être utilisé comme argument de table dans une fonction TVF.

scalar()

Renvoie un objet Column pour une sous-requête scalaire contenant exactement une ligne et une colonne.

exists()

Renvoie un objet Column pour une sous-requête EXISTS.

Méthode

Description

isLocal()

Renvoie `Vrai` si les méthodes `collect` et `take` peuvent être exécutées localement.

asTable()

Convertit le DataFrame en un objet TableArg, qui peut être utilisé comme argument de table dans une fonction TVF.

scalar()

Renvoie un objet Column pour une sous-requête scalaire contenant exactement une ligne et une colonne.

exists()

Renvoie un objet Column pour une sous-requête EXISTS.

Exemples

Opérations de base sur les DataFrame

Python
# Create a DataFrame
people = spark.createDataFrame([
{"deptId": 1, "age": 40, "name": "Alice", "gender": "M", "salary": 50},
{"deptId": 1, "age": 50, "name": "Bob", "gender": "M", "salary": 100},
{"deptId": 2, "age": 60, "name": "Sue", "gender": "F", "salary": 150},
{"deptId": 3, "age": 20, "name": "Tom", "gender": "M", "salary": 200}
])

# Select columns
people.select("name", "age").show()

# Filter rows
people.filter(people.age > 30).show()

# Add a new column
people.withColumn("age_plus_10", people.age + 10).show()

Agrégation et regroupement

Python
# Group by and aggregate
people.groupBy("gender").agg({"salary": "avg", "age": "max"}).show()

# Multiple aggregations
from pyspark.sql import functions as F
people.groupBy("deptId").agg(
F.avg("salary").alias("avg_salary"),
F.max("age").alias("max_age")
).show()

Jointures

Python
# Create another DataFrame
department = spark.createDataFrame([
{"id": 1, "name": "PySpark"},
{"id": 2, "name": "ML"},
{"id": 3, "name": "Spark SQL"}
])

# Join DataFrames
people.join(department, people.deptId == department.id).show()

Transformations complexes

Python
# Chained operations
result = people.filter(people.age > 30) \\
.join(department, people.deptId == department.id) \\
.groupBy(department.name, "gender") \\
.agg({"salary": "avg", "age": "max"}) \\
.sort("max(age)")
result.show()