【发布时间】:2021-07-08 02:18:37
【问题描述】:
我在处理不良记录和文件 (CSV) 时遇到了一些问题。 这是我的 CSV 文件
+------+---+---+----+
| Name| ID|int|int2|
+------+---+---+----+
| Sohel| 1| 4| 33|
| Sohel| 2| 5| 56|
| Sohel| 3| 6| 576|
| Sohel| a| 7| 567|
|Sohel2| c| 7| 567|
+------+---+---+----+
我正在使用预定义架构读取此文件
schema = StructType([
StructField("Name",StringType(),True),
StructField("ID",IntegerType(),True),
StructField("int",IntegerType(),True),
StructField("int2",IntegerType(),True),
StructField("_corrupt_record", StringType(),True)
])
df = spark.read.csv('dbfs:/tmp/test_file/test_csv.csv', header=True, schema=schema,
columnNameOfCorruptRecord='_corrupt_record')
结果是
+------+----+---+----+---------------+
| Name| ID|int|int2|_corrupt_record|
+------+----+---+----+---------------+
| Sohel| 1| 4| 33| null|
| Sohel| 2| 5| 56| null|
| Sohel| 3| 6| 576| null|
| Sohel|null| 7| 567| Sohel,a,7,567|
|Sohel2|null| 7| 567| Sohel2,c,7,567|
+------+----+---+----+---------------+
它给了我预期的结果,但是问题从这里开始我只想访问那些“_corrupt_record”并制作一个新的df。 我确实在 df 中过滤了“_corrupt_record”,但它似乎原始 CSV 文件没有“_corrupt_record”列,这就是它给我错误的原因。
badRows = df.filter("_corrupt_record is Not Null").show()
错误消息
Error while reading file dbfs:/tmp/test_file/test_csv.csv.
Caused by: java.lang.IllegalArgumentException: _corrupt_record does not exist. Available: Name, ID, int, int2
我正在流动 Databricks 文档, https://docs.databricks.com/data/data-sources/read-csv.html#read-files ,但是他们也有同样的错误,为什么他们甚至将它添加到文档中!!
我只想访问“_corrupt_record”列并制作新的 DF。 任何帮助或建议将不胜感激。
【问题讨论】:
标签: python pyspark databricks