代码之家  ›  专栏  ›  技术社区  ›  Gayatri

PySpark 2-合并来自多行的记录

  •  0
  • Gayatri  · 技术社区  · 8 年前

    我有一个文本文件,其记录如下:

    <BR>Datetime:2018.06.30^
    Name:ABC^
    Se:4^
    Machine:XXXXXXX^
    InnerTrace:^
    AdditionalInfo:^
    <ER>
    <BR>Datetime:2018.05.30-EDT^
    Name:DEF^
    Se:4^
    Machine:XXXXXXX^
    InnerTrace:^
    AdditionalInfo:^
    <ER>
    

    我正试图将其读入spark并处理该文件,以便得到以下结果:

    Datetime     Name  Se  Machine  InnerTrace  AdditionalInfo
    2018.06.30   ABC   4   XXXXXXX      
    2018.05.30   DEF   4   XXXXXXX
    

    当我尝试用

    sparkSession.read.csv("filename") 
    

    我将每一行分开,这使得很难将<BR>和<ER>之间的所有行放在一起。有什么简单的解决办法吗?

    我用PySpark 2做这个。

    1 回复  |  直到 8 年前
        1
  •  2
  •   pault Tanjin    8 年前

    这种文件格式不适合spark。如果你不能修改文件,你必须做大量的处理才能得到你想要的。

    以下是一种可能对你有用的方法-

    读取文件

    假设您有以下数据帧:

    df = spark.read.csv(path="filename", quote='')
    df.show(truncate=False)
    #+-----------------------------+
    #|_c0                          |
    #+-----------------------------+
    #|"<BR>Datetime:2018.06.30^    |
    #|Name:ABC^                    |
    #|Se:4^                        |
    #|Machine:XXXXXXX^             |
    #|InnerTrace:^                 |
    #|AdditionalInfo:^             |
    #|<ER>"                        |
    #|"<BR>Datetime:2018.05.30-EDT^|
    #|Name:DEF^                    |
    #|Se:4^                        |
    #|Machine:XXXXXXX^             |
    #|InnerTrace:^                 |
    #|AdditionalInfo:^             |
    #|<ER>"                        |
    #+-----------------------------+
    

    添加一列以将行分隔成记录组

    import pyspark.sql.functions as f
    from pyspark.sql import Window
    
    w = Window.orderBy("id").rangeBetween(Window.unboundedPreceding, 0)
    df = df.withColumn("group", f.col("_c0").rlike('^"<BR>.+').cast("int"))
    df = df.withColumn("id", f.monotonically_increasing_id())
    df = df.withColumn("group", f.sum("group").over(w)).drop("id")
    df.show(truncate=False)
    #+-----------------------------+-----+
    #|_c0                          |group|
    #+-----------------------------+-----+
    #|"<BR>Datetime:2018.06.30^    |1    |
    #|Name:ABC^                    |1    |
    #|Se:4^                        |1    |
    #|Machine:XXXXXXX^             |1    |
    #|InnerTrace:^                 |1    |
    #|AdditionalInfo:^             |1    |
    #|<ER>"                        |1    |
    #|"<BR>Datetime:2018.05.30-EDT^|2    |
    #|Name:DEF^                    |2    |
    #|Se:4^                        |2    |
    #|Machine:XXXXXXX^             |2    |
    #|InnerTrace:^                 |2    |
    #|AdditionalInfo:^             |2    |
    #|<ER>"                        |2    |
    #+-----------------------------+-----+
    

    使用Regex清理和拆分字符串

    df = df.select(
        "group",
        f.regexp_replace(pattern=r'(^"<BR>|<ER>"$|\^$)', replacement='', str="_c0").alias("col")
    ).where(f.col("col") != '')
    df = df.select("group", f.split("col", ":").alias("split"))
    
    df.show(truncate=False)
    #+-----+--------------------------+
    #|group|split                     |
    #+-----+--------------------------+
    #|1    |[Datetime, 2018.06.30]    |
    #|1    |[Name, ABC]               |
    #|1    |[Se, 4]                   |
    #|1    |[Machine, XXXXXXX]        |
    #|1    |[InnerTrace, ]            |
    #|1    |[AdditionalInfo, ]        |
    #|2    |[Datetime, 2018.05.30-EDT]|
    #|2    |[Name, DEF]               |
    #|2    |[Se, 4]                   |
    #|2    |[Machine, XXXXXXX]        |
    #|2    |[InnerTrace, ]            |
    #|2    |[AdditionalInfo, ]        |
    #+-----+--------------------------+
    

    从数组、groupby和pivot中提取元素

    df = df.select(
            "group",
            f.col("split").getItem(0).alias("key"),
            f.col("split").getItem(1).alias("value")
        )\
        .groupBy("group").pivot("key").agg(f.first("value"))\
        .drop("group")
    df.show(truncate=False)
    #+--------------+--------------+----------+-------+----+---+
    #|AdditionalInfo|Datetime      |InnerTrace|Machine|Name|Se |
    #+--------------+--------------+----------+-------+----+---+
    #|              |2018.06.30    |          |XXXXXXX|ABC |4  |
    #|              |2018.05.30-EDT|          |XXXXXXX|DEF |4  |
    #+--------------+--------------+----------+-------+----+---+
    
    推荐文章