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.
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 :
- Vous ne pouvez pas spécifier les options de la source de données.
- Vous ne pouvez pas spécifier le schéma pour les données.
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
- Scala
- R
# 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()
// Write wanderbricks reviews to CSV format
val reviews = spark.read.table("samples.wanderbricks.reviews")
reviews.write.format("csv").option("header", "true").save("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
// Read the CSV file into a DataFrame
val df = spark.read
.format("csv")
.option("header", "true")
.option("inferSchema", "true")
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
df.show()
df.printSchema()
df <- read.df("/Volumes/<catalog>/<schema>/<volume>/reviews_csv", source = "csv", header = "true", inferSchema = "true")
display(df)
printSchema(df)
Lire les fichiers CSV à l'aide de SQL
L'exemple SQL suivant lit un fichier CSV en utilisant read_files.
-- 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
- Scala
- SQL
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()
import org.apache.spark.sql.types._
val schema = StructType(Array(
StructField("review_id", StringType, nullable = true),
StructField("rating", IntegerType, nullable = true),
StructField("comment", StringType, nullable = true)
))
val df = spark.read.format("csv").schema(schema).option("header", "true").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
df.printSchema()
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
format => 'csv',
header => true,
schema => 'review_id string, rating int, comment string'
)
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
- Scala
- SQL
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)
import org.apache.spark.sql.types._
val schema = StructType(Array(
StructField("review_id", StringType, nullable = true),
StructField("rating", IntegerType, nullable = true)
))
val df = spark.read.format("csv").schema(schema).option("header", "true").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
df.show()
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
format => 'csv',
header => true,
schema => 'review_id string, rating int'
)
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 correctementDROPMALFORMED: supprime les lignes contenant des champs qui n'ont pas pu être analysésFAILFASTAnnule la lecture si des données malformées sont trouvées.
Pour définir le mode, utilisez l'option mode.
- Python
- Scala
- SQL
df = (spark.read
.format("csv")
.option("header", "true")
.option("mode", "PERMISSIVE")
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
)
val df = spark.read
.format("csv")
.option("header", "true")
.option("mode", "PERMISSIVE")
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
format => 'csv',
header => true,
mode => 'PERMISSIVE'
)
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
badRecordsPathpour enregistrer les enregistrements corrompus dans un fichier. - Vous pouvez ajouter la colonne
_corrupt_recordau schéma fourni au DataFrameReader pour examiner les enregistrements corrompus dans le DataFrame résultant.
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
- Scala
- SQL
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()))
import org.apache.spark.sql.types._
val schema = StructType(Array(
StructField("review_id", StringType, nullable = true),
StructField("rating", IntegerType, nullable = true),
StructField("comment", StringType, nullable = true),
StructField("_corrupt_record", StringType, nullable = true)
))
val df = spark.read
.format("csv")
.option("header", "true")
.option("mode", "PERMISSIVE")
.schema(schema)
.load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
df.filter(df("_corrupt_record").isNotNull).show()
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
format => 'csv',
header => true,
mode => 'PERMISSIVE',
schema => 'review_id string, rating int, comment string, _corrupt_record string'
)
WHERE _corrupt_record IS NOT NULL
Activez la colonne de données sauvée
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
- Scala
- SQL
df = spark.read.option("rescuedDataColumn", "_rescued_data").format("csv").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
val df = spark.read.option("rescuedDataColumn", "_rescued_data").format("csv").load("/Volumes/<catalog>/<schema>/<volume>/reviews_csv")
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
format => 'csv',
header => true,
rescuedDataColumn => '_rescued_data'
)
Pour supprimer le chemin du fichier source de la colonne de données récupérées, définissez :
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_recordoubadRecordsPath.
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.