CSV-fájlok olvasása és írása

A CSV (vesszővel tagolt értékek) egyszerű szöveges táblázatos formátum, amelyet széles körben használnak adatcserére, ETL-folyamatokra és általános célú adattárolásra. Azure Databricks támogatja a CSV-t az Apache Sparkkal való olvasáshoz és íráshoz, beleértve a sémakövetkeztetést, a tömörítést, a hibásan formázott rekordkezelést és a mentett adatokat.

Megjegyzés

A Databricks a CSV-fájlok olvasásához javasolja az read_files SQL-felhasználók számára a táblaértékű függvényt . read_files a Databricks Runtime 13.3 LTS-ben és újabb verziókban érhető el.

Ideiglenes nézetet is használhat. Ha az SQL-t használja a CSV-adatok közvetlen olvasására ideiglenes nézetek használata nélkül, vagy read_filesaz alábbi korlátozások érvényesek:

Prerequisites

Azure Databricks nem igényel további konfigurációt a CSV-fájlok használatához. A CSV-fájlok streameléséhez azonban automatikus betöltőre van szükség.

Beállítások

A CSV-adatforrások konfigurálásához használja a(z) .option().options()DataFrameReader és DataFrameWriter metódusait. A támogatott beállítások teljes listájáért tekintse meg DataFrameReader a CSV-beállításokat és DataFrameWriter a CSV-beállításokat.

Usage

Az alábbi példák bemutatják a CSV-fájlok olvasását és írását, a sémák megadását és a hibásan formázott rekordok kezelését.

CSV-fájlok olvasása

Az alábbi példa a Wanderbricks-mintaadatkészletet használja. Kiírja az értékelési adatokat CSV-be, majd visszaolvassa őket.

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-fájlok olvasása AZ SQL használatával

Az alábbi SQL-példa egy CSV-fájlt olvas be a használatával 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')

Séma megadása

Ha a CSV-fájl sémája ismert, megadhatja a kívánt sémát a CSV-olvasónak a schema beállítással.

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

Az oszlopok egy részhalmazának olvasása

A CSV-elemző viselkedése attól függ, hogy mely oszlopok lesznek beolvasva. Ha a megadott séma nem egyezik a fájlelrendezéssel, az eredmények jelentősen eltérhetnek attól függően, hogy mely oszlopok érhetők el. A CSV nem rendelkezik oszlopnév-metaadatokkal, ezért a Spark a sémamezőket oszlopokhoz rendeli pozíció szerint – a nem egyező séma nem megfelelő mezőkbe helyezi át az értékeket.

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

Helytelenül formázott CSV-rekordok kezelése

Ha csv-fájlokat olvas egy megadott sémával, lehetséges, hogy a fájlok adatai nem egyeznek a sémával. A város nevét tartalmazó mező például nem egész számként értelmezi. A következmények attól függenek, hogy az elemző milyen módon fut:

  • PERMISSIVE (alapértelmezett): a program null értékeket szúr be olyan mezőkhöz, amelyek nem elemezhetők megfelelően
  • DROPMALFORMED: nem elemezhető mezőket tartalmazó sorokat csepegtet
  • FAILFAST: megszakítja az olvasást, ha hibásan formázott adatok találhatók

A mód beállításához használja a mode lehetőséget.

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

Ebben a PERMISSIVE módban az alábbi módszerek egyikével megvizsgálhatók azok a sorok, amelyeket nem sikerült megfelelően elemezni:

  • Megadhat egy egyéni elérési utat a sérült rekordok fájlba történő rögzítésének lehetőségéhez badRecordsPath .
  • Az oszlopot _corrupt_record hozzáadhatja a DataFrameReaderhez biztosított sémához az eredményül kapott DataFrame sérült rekordjainak áttekintéséhez.

Megjegyzés

A badRecordsPath beállítás élvez elsőbbséget a _corrupt_record beállítással szemben, ami azt jelenti, hogy a megadott elérési útra írt hibásan formázott sorok nem jelennek meg az eredményül kapott DataFrame-ben.

A hibásan formázott rekordok alapértelmezett viselkedése megváltozik a mentett adatoszlop használatakor.

A hibás formátumú sorok _corrupt_record használatával történő vizsgálatához adja hozzá azt a sémához, majd szűrjön a nem null értékek alapján:

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

A helyreállított adatok oszlopának engedélyezése

Megjegyzés

Ez a funkció a Databricks Runtime 8.3-ban és újabb verziókban támogatott.

A mód használatakor PERMISSIVE engedélyezheti, hogy a mentett adatoszlop rögzítse azokat az adatokat, amelyeket nem elemeztek, mert egy rekord egy vagy több mezője az alábbi problémák egyikével jár:

  • Hiányzik a megadott sémából.
  • Nem egyezik a megadott séma adattípusával.
  • A megadott séma mezőneveivel kis- és nagybetű érzékenység tekintetében nem egyezik meg.

A mentett adatoszlop JSON-dokumentumként lesz visszaadva, amely a mentett oszlopokat és a rekord forrásfájl-elérési útját tartalmazza.

A helyreállított adatok oszlopának engedélyezéséhez beolvasáskor állítsa a(z) rescuedDataColumn beállítást egy oszlopnévre:

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

Ha el szeretné távolítani a forrásfájl elérési útját a mentett adatoszlopból, állítsa be a következőt:

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

A CSV-elemző három módot támogat a rekordok elemzésekor: PERMISSIVE, DROPMALFORMEDés FAILFAST. Ha rescuedDataColumn együtt használják, az adattípus eltérései nem okoznak rekordvesztést DROPMALFORMED módban, és nem jeleznek hibát FAILFAST módban. A rendszer csak a sérült rekordokat – azaz hiányos vagy hibásan formázott CSV-t – elveti vagy hibát jelez.

Ha rescuedDataColumn módban használjákPERMISSIVE, a sérült rekordokra a következő szabályok vonatkoznak:

  • A fájl első sora (fejlécsor vagy adatsor) beállítja a várt sorhosszt.
  • A különböző számú oszlopot tartalmazó sorok hiányosnak minősülnek.
  • Az adattípus eltérései nem tekinthetők sérült rekordnak.
  • A rendszer csak a hiányos és hibásan formázott CSV-rekordokat tekinti sérültnek, és ezeket a _corrupt_record oszlopba vagy a badRecordsPath oszlopba rögzíti.

További források

  • Parquet-fájlok olvasása és írása: Ha a számítási feladat jobb lekérdezési teljesítményt vagy hatékonyabb tárolást igényel, a Parquet oszlopos elrendezése jelentős előnyöket kínál a CSV egyszerű szöveges formátumával szemben.