【问题标题】:How to skip lines while reading a CSV file as a dataFrame using PySpark?如何在使用 PySpark 将 CSV 文件作为数据帧读取时跳过行?
【发布时间】:2017-05-25 00:27:27
【问题描述】:

我有一个这样构成的 CSV 文件:

Header
Blank Row
"Col1","Col2"
"1,200","1,456"
"2,000","3,450"

我在阅读这个文件时遇到了两个问题。

  1. 我想忽略标题并忽略空白行
  2. 值中的逗号不是分隔符

这是我尝试过的:

df = sc.textFile("myFile.csv")\
              .map(lambda line: line.split(","))\ #Split By comma
              .filter(lambda line: len(line) == 2).collect() #This helped me ignore the first two rows

但是,这不起作用,因为值中的逗号被读取为分隔符,而 len(line) 返回 4 而不是 2。

我尝试了另一种方法:

data = sc.textFile("myFile.csv")
headers = data.take(2) #First two rows to be skipped

当时的想法是使用过滤器而不是读取标题。但是,当我尝试打印标题时,我得到了编码值。

[\x00A\x00Y\x00 \x00J\x00u\x00l\x00y\x00 \x002\x000\x001\x006\x00]

读取 CSV 文件并跳过前两行的正确方法是什么?

【问题讨论】:

    标签: apache-spark pyspark spark-dataframe pyspark-sql


    【解决方案1】:

    尝试使用带有 'quotechar' 参数的 csv.reader。它将正确分割行。 之后,您可以根据需要添加过滤器。

    import csv
    from pyspark.sql.types import StringType
    
    df = sc.textFile("test2.csv")\
               .mapPartitions(lambda line: csv.reader(line,delimiter=',', quotechar='"')).filter(lambda line: len(line)>=2 and line[0]!= 'Col1')\
               .toDF(['Col1','Col2'])
    

    【讨论】:

    • 我通过调用csv.reader([l.replace('\0','') for l in line],delimiter=',', quotechar='"')修复了它
    【解决方案2】:

    对于您的第一个问题,只需使用zipWithIndex 压缩RDD 中的行并过滤您不想要的行。 对于第二个问题,您可以尝试从行中去除第一个和最后一个双引号字符,然后在 "," 上拆分该行。

    rdd = sc.textFile("myfile.csv")
    rdd.zipWithIndex().
        filter(lambda x: x[1] > 2).
        map(lambda x: x[0]).
        map(lambda x: x.strip('"').split('","')).
        toDF(["Col1", "Col2"])
    

    不过,如果您正在寻找在 Spark 中处理 CSV 文件的标准方法,最好使用 databricks 中的 spark-csv 包。

    【讨论】:

    • 为您的“虽然”投票 - 此外,该软件包不应与 Spark 2 一起使用,因为它已集成到 Spark 中,这使得“尽管”更加重要。我强烈建议在您的其他 Spark 逻辑之外的单独作业中进行这种过滤,因为这是经典的数据规范化/正则化,不应该成为分析管道的一部分。在 Spark 之外执行此操作可让您使用自定义工具来完成该工作,然后拥有每个人都可以使用的适当文件格式。
    【解决方案3】:

    Zlidime 的回答是正确的。工作解决方案是这样的:

    import csv
    
    customSchema = StructType([ \
        StructField("Col1", StringType(), True), \
        StructField("Col2", StringType(), True)])
    
    df = sc.textFile("file.csv")\
            .mapPartitions(lambda partition: csv.reader([line.replace('\0','') for line in partition],delimiter=',', quotechar='"')).filter(lambda line: len(line) > 2 and line[0] != 'Col1')\
            .toDF(customSchema)
    

    【讨论】:

      【解决方案4】:

      如果CSV文件结构总是有两列,在Scala上可以实现:

      val struct = StructType(
        StructField("firstCol", StringType, nullable = true) ::
        StructField("secondCol", StringType, nullable = true) :: Nil)
      
      val df = sqlContext.read
        .format("com.databricks.spark.csv")
        .option("header", "false")
        .option("inferSchema", "false")
        .option("delimiter", ",")
        .option("quote", "\"")
        .schema(struct)
        .load("myFile.csv")
      
      df.show(false)
      
      val indexed = df.withColumn("index", monotonicallyIncreasingId())
      val filtered = indexed.filter(col("index") > 2).drop("index")
      
      filtered.show(false)
      

      结果是:

      +---------+---------+
      |firstCol |secondCol|
      +---------+---------+
      |Header   |null     |
      |Blank Row|null     |
      |Col1     |Col2     |
      |1,200    |1,456    |
      |2,000    |3,450    |
      +---------+---------+
      
      +--------+---------+
      |firstCol|secondCol|
      +--------+---------+
      |1,200   |1,456    |
      |2,000   |3,450    |
      +--------+---------+
      

      【讨论】:

      • PySpark 允许你做同样的事情。如果不是标题,这将起作用。只有标题被读入,其他行被跳过。
      • 单调递增的 id 不能保证连续的 Id,所以如果你使用 >2 条件,你不能保证这会删除前 2 行
      【解决方案5】:

      您为什么不试试pyspark.sqlDataFrameReader API?这很容易。对于这个问题,我想这一行就足够了。

      df = spark.read.csv("myFile.csv") # By default, quote char is " and separator is ','
      

      使用此 API,您还可以使用一些其他参数,例如标题行,忽略前导和尾随空格。这是链接:DataFrameReader API

      【讨论】:

      • 这不允许我跳过行
      • 您是否尝试过将ignoreLeadingWhiteSpaceignoreTrailingWhiteSpace 设置为True?我不确定它会起作用,但至少,试一试。
      • 也可以试试mode=DROPMALFORMED。我的假设是它会认为空行已损坏。
      • 我试过 mode=DROPMALFORMED ,虽然它可能会丢弃空行,但它不会跳过标题(第一行)读取第一行后,预计只有一列。所以当它读取一行多于一列时会出错。
      猜你喜欢
      • 2018-02-14
      • 1970-01-01
      • 2020-07-20
      • 2020-02-04
      • 2013-09-24
      • 1970-01-01
      • 1970-01-01
      • 2018-11-14
      相关资源
      最近更新 更多