Чтение и запись CSV-файлов

CSV (значения с разделительной запятыми) — это формат табличной таблицы обычного текста, широко используемый для обмена данными, конвейеров ETL и хранилища данных общего назначения. Azure Databricks поддерживает CSV-файл для чтения и записи с помощью Apache Spark, включая вывод схемы, сжатие, обработку неправильно сформированных записей и спасенных данных.

Примечание.

Databricks рекомендует read_files функцию с табличным значением для пользователей SQL читать CSV-файлы. read_files доступен в Databricks Runtime 13.3 LTS и выше.

Можно также использовать временное представление. Если вы используете SQL для чтения данных CSV напрямую без использования временных представлений или read_files, применяются следующие ограничения:

Необходимые условия

Azure Databricks не требует дополнительной конфигурации для использования CSV-файлов. Однако для потоковой передачи CSV-файлов требуется автозагрузчик.

Параметры

Используйте методы .option() и .options() объектов DataFrameReader и DataFrameWriter для настройки источников данных CSV. Полный список поддерживаемых параметров см. в разделах DataFrameReaderпараметры CSV и DataFrameWriterпараметры CSV.

Usage

В следующих примерах показано чтение и запись CSV-файлов, указание схем и обработка неправильных записей.

Чтение CSV-файлов

В следующем примере используется пример набора данных Wanderbricks . Записывает данные отзывов в CSV-файл, а затем считывает их обратно.

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()

Scala

// 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()

R

df <- read.df("/Volumes/<catalog>/<schema>/<volume>/reviews_csv", source = "csv", header = "true", inferSchema = "true")
display(df)
printSchema(df)

Чтение CSV-файлов с помощью SQL

В следующем примере SQL считывается CSV-файл с помощью read_files.

-- mode "FAILFAST" aborts file parsing with a RuntimeException if malformed lines are encountered
SELECT * FROM read_files(
  'abfss://<bucket>@<storage-account>.dfs.core.windows.net/<path>/<file>.csv',
  format => 'csv',
  header => true,
  mode => 'FAILFAST')

Указание схемы

При известной схеме CSV-файла можно указать нужную схему для средства чтения CSV с помощью параметра 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()

Scala

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()

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
  format => 'csv',
  header => true,
  schema => 'review_id string, rating int, comment string'
)

Чтение подмножества столбцов

Поведение средства синтаксического анализа CSV зависит от того, какие столбцы считываются. Если указанная схема не соответствует макету файла, результаты могут значительно отличаться в зависимости от того, к каким столбцам обращаются. CSV-файл не содержит метаданных имени столбца, поэтому Spark сопоставляет поля схемы с столбцами по позиции— несогласованная схема перемещает значения в неправильные поля.

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)

Scala

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()

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
  format => 'csv',
  header => true,
  schema => 'review_id string, rating int'
)

Обработка неправильно сформированных записей CSV

При чтении CSV-файлов с указанной схемой возможно, что данные в файлах не соответствуют схеме. Например, поле, содержащее название города, не будет анализироваться как целое число. Последствия зависят от режима, в котором выполняется средство синтаксического анализа:

  • PERMISSIVE (по умолчанию): нулевые значения вставляются в поля, которые не удалось правильно проанализировать
  • DROPMALFORMED: удаляет строки, содержащие поля, которые не удалось проанализировать
  • FAILFAST: прерывает чтение, если найдены неправильные данные

Чтобы задать режим, используйте параметр mode.

Python

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

Scala

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

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
  format => 'csv',
  header => true,
  mode => 'PERMISSIVE'
)

В режиме PERMISSIVE строки, которые не удалось правильно проанализировать, можно проверить с помощью одного из следующих методов:

  • Можно указать пользовательский путь к параметру badRecordsPath для записи поврежденных записей в файл.
  • Вы можете добавить столбец _corrupt_record в схему, указанную в DataFrameReader, чтобы просмотреть поврежденные записи в результирующем кадре данных.

Примечание.

Параметр badRecordsPath имеет приоритет над _corrupt_record, то есть неправильные строки, записанные по указанному пути, не отображаются в результирующем кадре данных.

Поведение по умолчанию для неправильно сформированных записей изменяется при использовании восстановленного столбца данных.

Чтобы проверить неправильно сформированные строки с помощью _corrupt_record, добавьте его в схему и отфильтруйте значения, отличные от NULL:

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()))

Scala

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()

SQL

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

Включите столбец с восстановленными данными

Примечание.

Эта функция поддерживается в Databricks Runtime 8.3 и выше.

При использовании PERMISSIVE режима можно включить спасаемый столбец данных для записи любых данных, которые не были проанализированы, так как в одной или нескольких полях записи возникает одна из следующих проблем:

  • Отсутствует в предоставленной схеме.
  • Не соответствует типу данных предоставленной схемы.
  • Имеется несовпадение регистра с названиями полей в предоставленной схеме.

Спасательный столбец данных возвращается в виде документа JSON, содержащего столбцы, которые были спасены, и путь к исходному файлу записи.

Чтобы включить резервный столбец данных, задайте для параметра rescuedDataColumn имя столбца при чтении:

Python

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

Scala

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

SQL

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_csv',
  format => 'csv',
  header => true,
  rescuedDataColumn => '_rescued_data'
)

Чтобы удалить путь к исходному файлу из спасаемого столбца данных, установите следующий параметр:

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

Анализатор CSV поддерживает три режима анализа записей: PERMISSIVE, DROPMALFORMED и FAILFAST. При использовании вместе с rescuedDataColumn несоответствие типов данных не приводит к удалению записей в режиме DROPMALFORMED или возникновению ошибки в режиме FAILFAST. Только поврежденные записи, то есть неполные или некорректно сформированные CSV, будут удалены или вызовут ошибки.

При использовании rescuedDataColumn в режиме PERMISSIVE к поврежденным записям применяются следующие правила:

  • Первая строка файла (строка заголовка или строка данных) задает ожидаемую длину строки.
  • Строка с другим числом столбцов считается неполной.
  • Несоответствия типов данных не считается повреждением записей.
  • Только неполные и неправильно сформированные записи CSV считаются поврежденными и записываются в столбец _corrupt_record или badRecordsPath.

Дополнительные ресурсы

  • Чтение и запись файлов Parquet. Если рабочая нагрузка требует повышения производительности запросов или более эффективного хранилища, макет столбца Parquet предлагает значительные преимущества по сравнению с форматом обычного текста CSV.