【问题标题】:How do I read a non standard csv file into dataframe with python or scala如何使用 python 或 scala 将非标准 csv 文件读入数据帧
【发布时间】:2020-05-10 09:18:26
【问题描述】:

我在下面有一个数据集采样要使用 python 或 scala 处理:

FWD,13032009:09:01,10.56| FWD,13032009:10:53,11.23| FWD,13032009:15:40,23.20
SPOT,13032009:09:04,11.56| FWD,13032009:11:45,11.23| SPOT,13032009:12:30,23.20
FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20| FWD,13032009:14:340,56.00
FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20

每一行都将被分割成多个更小的字符串,可以进一步分割。

我正在寻找一种有效的方法来生成具有以下内容的 RDD 或 Dataframe:

FWD,13032009:09:01,10.56 
FWD,13032009:10:53,11.23
FWD,13032009:15:40,23.20
SPOT,13032009:09:04,11.56
FWD,13032009:11:45,11.23
SPOT,13032009:12:30,23.20
FWD,13032009:08:01,10.56
SPOT,13032009:12:30,11.23
FWD,13032009:13:20,23.20
FWD,13032009:14:340,56.00
FWD,13032009:08:01,10.56
SPOT,13032009:12:30,11.23
FWD,13032009:13:20,23.20

注意效率越高越好,因为生产中的总行数可能高达百万

非常感谢。

【问题讨论】:

    标签: python scala dataframe rdd


    【解决方案1】:

    假设您正在读取 csv 文件,您可以将每一行读取到一个列表中。展平这些值,然后将它们作为单独的行进行处理。

    将文件读入列表 - 100 万行不应过多处理:

    import csv
    import itertools
    
    import pandas as pd
    
    with open('test.csv','r') as f:
        reader = csv.reader(f, delimiter = '|')
        rows = list(reader)
    

    扁平化并从单个列表中拆分 - Python 标准库中出色的 itertools 库返回一个有助于内存且高效的生成器。

    flat_rows = itertools.chain.from_iterable(rows)
    list_rows = [i.strip().split(',') for i in flat_rows]
    

    嵌套列表list_rows 现在为您提供了一个干净且格式化的列表,如果您想创建dataframe,您可以将其发送至pandas

    list_rows
    >>
    [['FWD', '13032009:09:01', '10.56'],
     ['FWD', '13032009:10:53', '11.23'],
     ['FWD', '13032009:15:40', '23.20'],
     ['SPOT', '13032009:09:04', '11.56'],
     ['FWD', '13032009:11:45', '11.23'],
     ['SPOT', '13032009:12:30', '23.20'],
     ['FWD', '13032009:08:01', '10.56'],
     ['SPOT', '13032009:12:30', '11.23'],
     ['FWD', '13032009:13:20', '23.20'],
     ['FWD', '13032009:14:340', '56.00'],
     ['FWD', '13032009:08:01', '10.56'],
     ['SPOT', '13032009:12:30', '11.23'],
     ['FWD', '13032009:13:20', '23.20']]
    
    df = pd.DataFrame(list_rows)
    

    【讨论】:

    • 非常感谢 Bernard,我如何将 pandas df 注册到像常规 df 一样的临时表中,以便我可以使用 sql 进行分析?
    • @mdivk 好吧,我建议如果您想分析它,并且由于它已经位于数据框中,请使用pandas 进行分析。如果你想存储它然后查找pd.DataFrame.to_sql。如果您有新问题,请发布一个新问题并从那里开始,谢谢!
    • 谢谢伯纳德,我只需要将 df 注册到临时表中,然后运行 ​​ad-hoc sql 查询,你能在这里分享一下吗......请.....
    • 我试过 df.registerTempTable("test") 但它在我的 python3 Jupyter 笔记本中不起作用
    • 你需要清楚你是在spark rdd 还是pandas 数据帧,他们有不同的方法。
    【解决方案2】:

    Python 解决方案:如果您将文本作为字符串获取,您可以使用换行符 (\n) 对您的 序列进行 replace() 序列,然后将其作为 DataFrame 读取:

    import pandas as pd
    from io import StringIO
    
    data_set = """FWD,13032009:09:01,10.56| FWD,13032009:10:53,11.23| FWD,13032009:15:40,23.20
    SPOT,13032009:09:04,11.56| FWD,13032009:11:45,11.23| SPOT,13032009:12:30,23.20
    FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20| FWD,13032009:14:340,56.00
    FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20
    """
    data_set *= 100000  # Make it over a million elements to ensure performance is adequate
    data_set = data_set.replace("| ", "\n")
    
    data_set_stream = StringIO(data_set)  # Pandas needs to read a file-like object, so need to turn our string into a buffer
    df = pd.read_csv(data_set_stream)
    print(df)  # df is our desired DataFrame
    

    【讨论】:

    • 谢谢 Oliver,在生产中我们无法按照您的方式生成数据集,它必须从文件中读取,这里的问题是文件不是标准的 csv,因为有些行包含不同的较小的csv 部分比其他部分,在示例数据中,您可以看到第三行有 4 个部分,而其他部分有 3 个
    • 您可以简单地将文件读入data_set 变量:with open("DATASET_FILENAME", "r") as my_file: data_set = my_file.read(),而不是硬编码字符串(如我的示例中)。其次,我的示例通过分离发现要连接的任意数量的部分来巧妙地处理不同数量的部分(您可以尝试我的代码并查看它是否正确地展平了数据集)。
    • 谢谢 Oliver,如何将 pandas 创建的数据框注册为 temptable?我想在 teamptable 上运行查询
    • 如果您想运行查询,pandas 有一个 api here
    【解决方案3】:

    如果您有兴趣,这里是 Scala 方式,

    val rdd1 = sc.parallelize(List("FWD,13032009:09:01,10.56| FWD,13032009:10:53,11.23| FWD,13032009:15:40,23.20", "SPOT,13032009:09:04,11.56| FWD,13032009:11:45,11.23| SPOT,13032009:12:30,23.20","FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20| FWD,13032009:14:340,56.00","FWD,13032009:08:01,10.56| SPOT,13032009:12:30,11.23| FWD,13032009:13:20,23.20"))

    val rdd2 = rdd1.flatMap(l => l.replaceAll(" ","").split("\\|")) val rds = rdd2.toDS

    val df = spark.read.csv(rds)

    df.show(false) +----+---------------+-----+ |_c0 |_c1 |_c2 | +----+---------------+-----+ |FWD |13032009:09:01 |10.56| |FWD |13032009:10:53 |11.23| |FWD |13032009:15:40 |23.20| |SPOT|13032009:09:04 |11.56| |FWD |13032009:11:45 |11.23| |SPOT|13032009:12:30 |23.20| |FWD |13032009:08:01 |10.56| |SPOT|13032009:12:30 |11.23| |FWD |13032009:13:20 |23.20| |FWD |13032009:14:340|56.00| |FWD |13032009:08:01 |10.56| |SPOT|13032009:12:30 |11.23| |FWD |13032009:13:20 |23.20| +----+---------------+-----+

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-04-17
      • 1970-01-01
      • 1970-01-01
      • 2017-04-29
      • 2019-01-30
      • 2018-11-14
      • 1970-01-01
      • 2021-11-05
      相关资源
      最近更新 更多