【问题标题】:Apache Spark working with pipe delimited CSV filesApache Spark 使用管道分隔的 CSV 文件
【发布时间】:2019-04-28 09:12:32
【问题描述】:

我对 Apache Spark 非常陌生,我正在尝试将 SchemaRDD 与我的管道分隔文本文件一起使用。我在我的 Mac 上使用 Scala 10 独立安装了 Spark 1.5.2。我有一个包含以下代表性数据的 CSV 文件,我试图根据记录的第一个值(列)将以下文件拆分为 4 个不同的文件.我非常感谢我能得到的任何帮助。

1|1.8|20140801T081137|115810740
2|20140714T060000|335|22159892|3657|0.00|||181
2|20140714T061500|335|22159892|3657|0.00|||157
2|20140714T063000|335|22159892|3657|0.00|||156
2|20140714T064500|335|22159892|3657|0.00|||66
2|20140714T070000|335|22159892|3657|0.01|||633
2|20140714T071500|335|22159892|3657|0.01|||1087
3|34|Starz
3|35|VH1
3|36|CSPAN: Cable Satellite Public Affairs Network
3|37|Encore
3|278|CMT: Country Music Television
3|281|Telehit
4|625363|1852400|Matlock|9212|The Divorce
4|625719|1852400|Matlock|16|The Rat Pack
4|625849|1846952|Smallville|43|Calling

【问题讨论】:

  • 欢迎来到 SO。如果您包括自己的尝试,您将有更好的机会获得答案。

标签: scala apache-spark apache-spark-sql


【解决方案1】:

注意:您的 csv 文件在每行中没有相同数量的字段 - 这不能按原样解析为 DataFrame。 (SchemaRDD 已重命名为 DataFrame。)如果您的 csv 文件格式正确,您可以执行以下操作:

使用 --packages com.databricks:spark-csv_2.10:1.3.0 启动 spark-shell 或 spark-submit 以便轻松解析 csv 文件 (see here)。在 Scala 中,您的代码将是,假设您的 csv 文件有一个标题 - 如果是,则更容易引用列:

val df = sqlContext.read.format("com.databricks.spark.csv").option("header", "true").option("inferSchema", "true").option("delimiter", '|').load("/path/to/file.csv")
// assume 1st column has name col1
val df1 = df.filter( df("col1") === 1)  // 1st DataFrame
val df2 = df.filter( df("col1") === 2)  // 2nd DataFrame  etc... 

由于您的文件格式不正确,您必须以不同的方式解析每一行,例如,执行以下操作:

val lines = sc.textFile("/path/to/file.csv")

case class RowRecord1( col1:Int, col2:Double, col3:String, col4:Int)
def parseRowRecord1( arr:Array[String]) = RowRecord1( arr(0).toInt, arr(1).toDouble, arr(2), arr(3).toInt)

case class RowRecord2( col1:Int, col2:String, col3:Int, col4:Int, col5:Int, col6:Double, col7:Int)
def parseRowRecord2( arr:Array[String]) = RowRecord2( arr(0).toInt, arr(1), arr(2).toInt, arr(3).toInt, arr(4).toInt, arr(5).toDouble, arr(8).toInt)

val df1 = lines.filter(_.startsWith("1")).map( _.split('|')).map( arr => parseRowRecord1( arr )).toDF
val df2 = lines.filter(_.startsWith("2")).map( _.split('|')).map( arr => parseRowRecord2( arr )).toDF

【讨论】:

  • 您好 KrisP,非常感谢您的帮助。我尝试了你的前几行代码,效果很好!我将尝试您的示例的其余部分,然后根据 COL0 的值将文件(具有不同的 num 列)拆分为具有相同列数的多个文件...
  • 嗨 KrisP,您还知道如何将输出保存到管道分隔文件吗?我认为 df2.write.format("com.databricks.spark.csv").save("/Users/temp/parsed1.txt") 命令的输出默认以逗号分隔,并将其分解为多个文件。如果可能,我还尝试将结果直接写入 Amazon Redshift 以使工作流程更加简化。非常感谢您的所有帮助。
  • 这仅适用于我使用双引号作为分隔符:.option("delimiter","|") 否则我会收到错误:java.lang.IllegalArgumentException: Delimiter cannot be more than一个字符
【解决方案2】:

在 PySpark 中,命令是:

df = spark.read.csv("filepath", sep="|")

【讨论】:

    猜你喜欢
    • 2018-05-16
    • 1970-01-01
    • 1970-01-01
    • 2020-09-25
    • 2023-03-15
    • 2020-07-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多