讀取與寫入 CSV 檔案

CSV(逗號分隔值)是一種廣泛用於資料交換、ETL 管線及通用資料儲存的純文字表格格式。 Azure Databricks 支援使用 Apache Spark 讀取及寫入 CSV,包括結構描述推斷、壓縮、格式錯誤記錄處理及救援資料。

注意

Databricks 建議 SQL 使用者讀取 CSV 檔案的數據表值函式。read_files Databricks Runtime 13.3 LTS 及以上版本中提供 read_files。

您也可以使用暫存檢視。 如果您使用 SQL 直接讀取 CSV 資料而不使用暫存檢視或 read_files,則適用下列限制:

先決條件

Azure Databricks 使用 CSV 檔案不需要額外設定。 不過,要串流 CSV 檔案,你需要 Auto Loader。

選項

使用 .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)

使用 SQL 讀取 CSV 檔案

下列 SQL 範例會使用 read_files 讀取 CSV 檔案。

-- 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 檔案

你也可以使用 USING CSV 子句建立以 CSV 檔案為基礎的暫存檢視,然後使用 SQL 進行查詢。 在子 OPTIONS 句中傳遞檔案路徑和讀取器的選項。 該視圖僅對當前會話可見,且在會話結束時會被移除。

CREATE TEMPORARY VIEW diamonds
USING CSV
OPTIONS (path "/databricks-datasets/Rdatasets/data-001/csv/ggplot2/diamonds.csv", header "true", mode "FAILFAST");

SELECT * FROM diamonds;

指定一個結構

當 CSV 檔案的結構描述為已知時,您可以使用 schema 選項,將所需的結構描述指定給 CSV 助讀程式。

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 的結構描述,藉此檢閱結果 DataFrame 中的損毀記錄。

注意

選項 badRecordsPath 的優先順序高於 _corrupt_record,這表示寫入所提供路徑的格式錯誤的資料列不會出現在產生的 DataFrame 中。

使用已修復的資料行時,格式錯誤的記錄的預設行為會變更。

若要使用 _corrupt_record 檢查格式錯誤的資料列,請將其加入結構描述中,並篩選非空值:

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 執行環境 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 的純文字格式有顯著優勢。