Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
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:
- Anda tidak dapat menentukan opsi sumber data.
- Anda tidak dapat menentukan skema untuk data.
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
badRecordsPathuntuk mencatat rekaman yang cacat ke file. - Anda dapat menambahkan kolom
_corrupt_recordke 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_recordataubadRecordsPath.
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.