Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
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_recordhozzá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_recordoszlopba vagy abadRecordsPathoszlopba 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.