Membaca dan menulis file CSV

CSV (nilai yang dipisahkan koma) adalah format tabular teks biasa yang banyak digunakan untuk pertukaran data, alur ETL, dan penyimpanan data tujuan umum. Azure Databricks mendukung CSV untuk membaca dan menulis dengan Apache Spark, termasuk inferensi skema, kompresi, penanganan rekaman cacat, dan data yang diselamatkan.

Catatan

Databricks read_files merekomendasikan fungsi bernilai tabel bagi pengguna SQL untuk membaca file CSV. read_files tersedia di Databricks Runtime 13.3 LTS ke atas.

Anda juga dapat menggunakan tampilan sementara. Jika Anda menggunakan SQL untuk membaca data CSV secara langsung tanpa menggunakan tampilan sementara atau read_files, batasan berikut berlaku:

Prerequisites

Azure Databricks tidak memerlukan konfigurasi tambahan untuk menggunakan file CSV. Namun, untuk melakukan streaming file CSV, Anda memerlukan Auto Loader.

Opsi

.option() Gunakan metode .options() dan DataFrameReader dan DataFrameWriter untuk mengonfigurasi sumber data CSV. Untuk daftar lengkap opsi yang didukung, lihat DataFrameReader Opsi CSV dan DataFrameWriter opsi CSV.

Usage

Contoh berikut menunjukkan membaca dan menulis file CSV, menentukan skema, dan menangani rekaman cacat.

Membaca file CSV

Contoh berikut menggunakan himpunan data sampel Wanderbricks . Ini menyimpan data ulasan ke file CSV, lalu membacanya kembali.

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)

Membaca file CSV menggunakan SQL

Contoh SQL berikut membaca file CSV menggunakan 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')

Tentukan skema

Ketika skema file CSV diketahui, Anda dapat menentukan skema yang diinginkan ke pembaca CSV dengan opsi 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'
)

Baca sebagian dari kolom

Perilaku pengurai CSV bergantung pada kolom mana yang dibaca. Jika skema yang ditentukan tidak cocok dengan tata letak file, hasilnya dapat sangat berbeda tergantung pada kolom mana yang diakses. CSV tidak memiliki metadata nama kolom, sehingga Spark memetakan bidang skema ke kolom menurut posisi — skema yang tidak cocok menggeser nilai ke bidang yang salah.

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

Menangani rekaman CSV yang salah bentuk

Saat membaca file CSV dengan skema yang ditentukan, ada kemungkinan bahwa data dalam file tidak cocok dengan skema. Misalnya, bidang yang berisi nama kota tidak akan mengurai sebagai bilangan bulat. Konsekuensinya tergantung pada mode yang dijalankan parser:

  • PERMISSIVE (default): null disisipkan untuk bidang yang tidak dapat diurai dengan benar
  • DROPMALFORMED: menghapus baris yang berisi kolom yang tidak dapat diparse
  • FAILFAST: membatalkan pembacaan jika ada data yang salah ditemukan

Untuk mengatur mode, gunakan opsi 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'
)

Dalam mode PERMISSIVE, dimungkinkan untuk memeriksa baris yang tidak dapat diurai secara tepat dengan salah satu metode berikut ini:

  • Anda dapat menyediakan jalur kustom ke opsi badRecordsPath untuk mencatat rekaman yang cacat ke file.
  • Anda dapat menambahkan kolom _corrupt_record ke skema yang disediakan ke DataFrameReader untuk meninjau rekaman yang rusak di DataFrame yang dihasilkan.

Catatan

Opsi badRecordsPath mendahului _corrupt_record, artinya baris yang terbentuk salah dan ditulis pada jalur yang disediakan tidak muncul dalam DataFrame hasil.

Perilaku default untuk rekaman rusak berubah saat menggunakan kolom data yang dipulihkan .

Untuk memeriksa baris cacat menggunakan _corrupt_record, tambahkan ke skema dan filter pada nilai non-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

Mengaktifkan kolom data yang diselamatkan

Catatan

Fitur ini didukung di Databricks Runtime 8.3 ke atas.

Saat menggunakan PERMISSIVE mode , Anda dapat mengaktifkan kolom data yang diselamatkan untuk mengambil data apa pun yang tidak diurai karena satu atau beberapa bidang dalam rekaman memiliki salah satu masalah berikut:

  • Tidak ada dalam skema yang disediakan.
  • Tidak cocok dengan jenis data skema yang disediakan.
  • Memiliki ketidakcocokan kasus dengan nama bidang dalam skema yang disediakan.

Kolom data yang diselamatkan dikembalikan sebagai dokumen JSON yang berisi kolom yang diselamatkan, dan jalur file sumber rekaman.

Untuk mengaktifkan kolom data yang diselamatkan, atur opsi rescuedDataColumn menjadi nama kolom saat membaca:

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

Untuk menghapus jalur file sumber dari kolom data yang diselamatkan, atur:

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

Pengurai CSV mendukung tiga mode saat mengurai baris: PERMISSIVE, DROPMALFORMED, dan FAILFAST. Ketika digunakan bersama dengan rescuedDataColumn, ketidakcocokan tipe data tidak menyebabkan penghapusan catatan dalam mode DROPMALFORMED atau menghasilkan kesalahan dalam mode FAILFAST. Hanya catatan yang rusak—yaitu CSV yang tidak lengkap atau cacat—yang dihapus atau memunculkan kesalahan.

Saat rescuedDataColumn digunakan pada mode PERMISSIVE, aturan berikut ini berlaku pada rekaman yang cacat:

  • Baris pertama file (baik baris header atau baris data) akan mengatur panjang baris yang diharapkan.
  • Baris dengan jumlah kolom yang berbeda dianggap tidak lengkap.
  • Ketidakcocokan jenis data tidak akan dianggap sebagai rekaman yang rusak.
  • Hanya rekaman CSV yang tidak lengkap dan salah format yang dianggap rusak dan direkam ke kolom _corrupt_record atau badRecordsPath.

Sumber daya tambahan

  • Membaca dan menulis file Parquet: Jika beban kerja Anda memerlukan performa kueri yang lebih baik atau penyimpanan yang lebih efisien, tata letak kolom Parquet menawarkan keuntungan signifikan daripada format teks biasa CSV.