Aller au contenu principal

Lisez et écrivez des fichiers CSV

Le CSV (valeurs séparées par des virgules) est un format tabulaire en texte brut largement utilisé pour l'échange de données, les pipelines ETL et le stockage de données à usage général. Databricks prend en charge le CSV pour la lecture et l'écriture avec Apache Spark, y compris l'inférence de schéma, la compression, la gestion des enregistrements mal formés et les données récupérées.

remarque

Databricks recommande la read_files fonction à valeur de table aux utilisateurs SQL pour lire les fichiers CSV. read_files est disponible dans Databricks Runtime 13.3 LTS et versions ultérieures.

Vous pouvez également utiliser une vue temporaire. Si vous utilisez SQL pour lire les données CSV directement sans utiliser de vues temporaires ou read_files, les limitations suivantes s'appliquent :

Prérequis

Databricks ne nécessite pas de configuration supplémentaire pour utiliser les fichiers CSV. Cependant, pour diffuser en continu des fichiers CSV, vous avez besoin d'Auto Loader.

Options

Utilisez les méthodes .option() et .options() de DataFrameReader et DataFrameWriter pour configurer les sources de données CSV. Pour obtenir la liste complète des options prises en charge, consultez DataFrameReader options CSV et DataFrameWriter options CSV.

Utilisation

Les exemples suivants illustrent la lecture et l'écriture de fichiers CSV, la spécification de schémas et la gestion des enregistrements mal formés.

Lire les fichiers CSV

L'exemple suivant utilise le dataset échantillon Wanderbricks. Il écrit des données de révision en CSV, puis les relit.

Python
# Write wanderbricks reviews to CSV format
df = spark.read.table("samples.wanderbricks.reviews")
df.write.format("csv").option("header", "true").save("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")

# Read the CSV file into a DataFrame
df = (spark.read
.format("csv")
.option("header", "true")
.option("inferSchema", "true")
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv"))
display(df)
df.printSchema()

Lire les fichiers CSV à l'aide de SQL

L'exemple SQL suivant lit un fichier CSV en utilisant read_files.

SQL
-- mode "FAILFAST" aborts file parsing with a RuntimeException if malformed lines are encountered
SELECT * FROM read_files(
's3://<bucket>/<path>/<file>.csv',
format => 'csv',
header => true,
mode => 'FAILFAST')

Spécifiez un schéma

Lorsque le schéma du fichier CSV est connu, vous pouvez spécifier le schéma souhaité au lecteur CSV avec l'option schema.

Python
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

schema = StructType([
StructField("review_id", StringType(), True),
StructField("rating", IntegerType(), True),
StructField("comment", StringType(), True)
])

df = spark.read.format("csv").schema(schema).option("header", "true").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
df.printSchema()

Lire un sous-ensemble de colonnes

Le comportement de l'analyseur CSV dépend des colonnes lues. Si le schéma spécifié ne correspond pas au Layout du fichier, les résultats peuvent varier considérablement selon les colonnes auxquelles on accède. Le fichier CSV ne contient pas de métadonnées de noms de colonne. Spark mappe donc les champs de schéma aux colonnes par position, et un schéma non concordant déplace les valeurs dans les mauvais champs.

Python
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

# Read only a subset of columns by specifying a partial schema
schema = StructType([
StructField("review_id", StringType(), True),
StructField("rating", IntegerType(), True)
])

df = spark.read.format("csv").schema(schema).option("header", "true").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
display(df)

Gérer les enregistrements CSV mal formés

Lors de la lecture de fichiers CSV avec un schéma spécifié, il est possible que les données des fichiers ne correspondent pas au schéma. Par exemple, un champ contenant le nom de la ville ne sera pas analysé comme un entier. Les conséquences dépendent du mode d'exécution de l'analyseur :

  • PERMISSIVE (default) : des valeurs null sont insérées pour les champs qui n'ont pas pu être analysés correctement
  • DROPMALFORMED: supprime les lignes contenant des champs qui n'ont pas pu être analysés
  • FAILFASTAnnule la lecture si des données malformées sont trouvées.

Pour définir le mode, utilisez l'option mode.

Python
df = (spark.read
.format("csv")
.option("header", "true")
.option("mode", "PERMISSIVE")
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
)

En mode PERMISSIVE, il est possible d'inspecter les lignes qui n'ont pas pu être analysées correctement en utilisant l'une des méthodes suivantes :

  • Vous pouvez fournir un chemin personnalisé à l'option badRecordsPath pour enregistrer les enregistrements corrompus dans un fichier.
  • Vous pouvez ajouter la colonne _corrupt_record au schéma fourni au DataFrameReader pour examiner les enregistrements corrompus dans le DataFrame résultant.
remarque

L'option badRecordsPath a la priorité sur _corrupt_record, ce qui signifie que les lignes mal formées écrites dans le chemin fourni n'apparaissent pas dans le DataFrame résultant.

Le comportement par default pour les enregistrements mal formés change lors de l'utilisation de la colonne de données sauvées.

Pour inspecter les lignes mal formées à l'aide de _corrupt_record, ajoutez-le au schéma et filtrez sur les valeurs non nulles :

Python
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

schema = StructType([
StructField("review_id", StringType(), True),
StructField("rating", IntegerType(), True),
StructField("comment", StringType(), True),
StructField("_corrupt_record", StringType(), True)
])

df = (spark.read
.format("csv")
.option("header", "true")
.option("mode", "PERMISSIVE")
.schema(schema)
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
)
display(df.filter(df["_corrupt_record"].isNotNull()))

Activez la colonne de données sauvée

remarque

Cette fonctionnalité est prise en charge dans Databricks Runtime 8.3 ou version supérieure.

Lorsque vous utilisez le mode PERMISSIVE, vous pouvez activer la colonne de données récupérées pour capturer toutes les données qui n'ont pas été analysées parce qu'un ou plusieurs champs d'un enregistrement présentent l'un des problèmes suivants :

  • Absent du schéma fourni.
  • Ne correspond pas au type de données du schéma fourni.
  • Il y a une incohérence de casse avec les noms des champs dans le schéma fourni.

La colonne des données sauvées est renvoyée sous forme de document JSON contenant les colonnes qui ont été sauvées, et le chemin de fichier source de l'enregistrement.

Pour activer la colonne de données récupérées, définissez l'option rescuedDataColumn sur un nom de colonne lors de la lecture :

Python
df = spark.read.option("rescuedDataColumn", "_rescued_data").format("csv").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")

Pour supprimer le chemin du fichier source de la colonne de données récupérées, définissez :

Python
spark.conf.set("spark.databricks.sql.rescuedDataColumn.filePath.enabled", "false")

L'analyseur CSV prend en charge trois modes lors de l'analyse des enregistrements : PERMISSIVE, DROPMALFORMED et FAILFAST. Lorsqu'elles sont utilisées avec rescuedDataColumn, les non-correspondances de type de données n'entraînent pas la suppression d'enregistrements en mode DROPMALFORMED ou le déclenchement d'une erreur en mode FAILFAST. Seuls les enregistrements corrompus, c'est-à-dire les fichiers CSV incomplets ou mal formés, sont supprimés ou génèrent des erreurs.

Lorsque rescuedDataColumn est utilisé en mode PERMISSIVE, les règles suivantes s'appliquent aux enregistrements corrompus:

  • La première ligne du fichier (qu'il s'agisse d'une ligne d'en-tête ou d'une ligne de données) définit la longueur de ligne attendue.
  • Une ligne avec un nombre de colonnes différent est considérée comme incomplète.
  • Les incompatibilités de type de données ne sont pas considérées comme des enregistrements corrompus.
  • Seuls les enregistrements CSV incomplets et mal formés sont considérés comme corrompus et enregistrés dans la colonne _corrupt_record ou badRecordsPath.

Ressources supplémentaires

  • Lire et écrire des fichiers Parquet: Si votre charge de travail exige une meilleure performance de query ou un stockage plus efficace, la disposition en colonnes de Parquet offre des avantages significatifs par rapport au format texte brut de CSV.