【问题标题】:How to use double pipe as delimiter in CSV?如何在 CSV 中使用双管道作为分隔符?
【发布时间】:2016-12-21 17:05:33
【问题描述】:

Spark 1.5 和 Scala 2.10.6

我有一个使用“¦¦”作为分隔符的数据文件。我很难通过这个解析来创建一个数据框。可以使用多个分隔符来创建一个数据框吗?该代码适用于单个损坏的管道,但不适用于多个分隔符。

我的代码:

val customSchema_1 = StructType(Array(
    StructField("ID", StringType, true), 
    StructField("FILLER", StringType, true), 
    StructField("CODE", StringType, true)));

val df_1 = sqlContext.read
    .format("com.databricks.spark.csv")
    .schema(customSchema_1)
    .option("delimiter", "¦¦")
    .load("example.txt")

示例文件:

12345¦¦  ¦¦10

【问题讨论】:

  • 你试过这个 ("\\|\\|") 吗?请see
  • 我建议将其转换为如下val text = sc.textFile("yourcsv.csv") val words = text.map( lines => lines.split("\\|\\|") 然后再次用单管道构造 csv 并继续您的方法
  • @RamGhadiyaram 如果 OP 的数据包含任何双管道可能会出现问题,我会尝试在 spark.csv 分隔符选项上使用转义。
  • @RamGhadiyaram 感谢您的建议!我厌倦了 .option("delimiter", "\\\\\\\\\\¦") 并且出现不支持的特殊字符错误。
  • @RamGhadiyaram ¦| 不同

标签: scala apache-spark


【解决方案1】:

我遇到了这个问题并找到了一个很好的解决方案,我使用的是 spark 2.3,我觉得它应该适用于 spark 2.2+ 的所有版本,但尚未对其进行测试。它的工作方式是我用 tab 替换 || ,然后内置的 csv 可以采用 Dataset[String] 。我使用制表符是因为我的数据中有逗号。

var df = spark.sqlContext.read
  .option("header", "true")
  .option("inferSchema", "true")
  .option("delimiter", "\t")
  .csv(spark.sqlContext.read.textFile("filename")
      .map(line => line.split("\\|\\|").mkString("\t")))

希望这对其他人有所帮助。

编辑:

从 spark 3.0.1 开始,这是开箱即用的。

示例:

val ds = List("name||id", "foo||12", "brian||34", """"cray||name"||123""", "cray||name||123").toDS
ds: org.apache.spark.sql.Dataset[String] = [value: string]

val csv = spark.read.option("header", "true").option("inferSchema", "true").option("delimiter", "||").csv(ds)
csv: org.apache.spark.sql.DataFrame = [name: string, id: string]

csv.show
+----------+----+
|      name|  id|
+----------+----+
|       foo|  12|
|     brian|  34|
|cray||name| 123|
|      cray|name|
+----------+----+

【讨论】:

  • 谢谢,这很完美(应该有更多的支持)!
【解决方案2】:

所以这里发出的实际错误是:

java.lang.IllegalArgumentException: Delimiter cannot be more than one character: ¦¦

文档证实了这一限制,我检查了 Spark 2.0 csv 阅读器,它具有相同的要求。

鉴于所有这些,如果您的数据足够简单,您不会有包含 ¦¦ 的条目,我会像这样加载您的数据:

scala> :pa
// Entering paste mode (ctrl-D to finish)
val customSchema_1 = StructType(Array(
    StructField("ID", StringType, true), 
    StructField("FILLER", StringType, true), 
    StructField("CODE", StringType, true)));

// Exiting paste mode, now interpreting.
customSchema_1: org.apache.spark.sql.types.StructType = StructType(StructField(ID,StringType,true), StructField(FILLER,StringType,true), StructField(CODE,StringType,true))

scala> val rawData = sc.textFile("example.txt")
rawData: org.apache.spark.rdd.RDD[String] = example.txt MapPartitionsRDD[1] at textFile at <console>:31

scala> import org.apache.spark.sql.Row
import org.apache.spark.sql.Row

scala> val rowRDD = rawData.map(line => Row.fromSeq(line.split("¦¦")))
rowRDD: org.apache.spark.rdd.RDD[org.apache.spark.sql.Row] = MapPartitionsRDD[3] at map at <console>:34

scala> val df = sqlContext.createDataFrame(rowRDD, customSchema_1)
df: org.apache.spark.sql.DataFrame = [ID: string, FILLER: string, CODE: string]

scala> df.show
+-----+------+----+
|   ID|FILLER|CODE|
+-----+------+----+
|12345|      |  10|
+-----+------+----+

【讨论】:

  • 如何添加 |^|在 spark 2 中另存为 csv 文件时的分隔符
  • 有时我们需要在不知道列名的情况下加载数据。我认为在这种情况下上述方法将失败
【解决方案3】:

我们尝试通过以下方式读取具有自定义分隔符和自定义数据框列名的数据,

# Hold new column names saparately
headers ="JC_^!~_*>Year_^!~_*>Date_^!~_*>Service_Type^!~_*>KMs_Run^!~_*>

# '^!~_*>' This is field delimiter, so split string
head = headers.split("^!~_*>")

## Below command splits the S3 file with custom delimiter and converts into Dataframe
df = sc.textFile("s3://S3_Path/sample.txt").map(lambda x: x.split("^!~_*>")).toDF(head)

在 toDF() 中将 head 作为参数传递给从具有自定义分隔符的文本文件创建的数据框。

希望这会有所帮助。

【讨论】:

    【解决方案4】:

    从 Spark2.8 及以上版本开始支持多字符分隔符。 https://issues.apache.org/jira/browse/SPARK-24540

    @lockwobr 提出的上述解决方案适用于 scala。在 Spark 2.8 以下工作并在 PySpark 中寻找解决方案的人可以参考以下内容

    ratings_schema = StructType([
                                      StructField("user_id", StringType(), False)
                                    , StructField("movie_id", StringType(), False)
                                    , StructField("rating", StringType(), False)
                                    , StructField("rating_timestamp", StringType(), True)
                                    ])
    
        #movies_df = spark.read.csv("ratings.dat", header=False, sep="::", schema=ratings_schema)
    
        movies_df = spark.createDataFrame(
                spark.read.text("ratings.dat").rdd.map(lambda line: line[0].split("::")),
                ratings_schema)

    我提供了一个示例,但您可以根据自己的逻辑对其进行修改。

    【讨论】:

      猜你喜欢
      • 2018-08-21
      • 1970-01-01
      • 1970-01-01
      • 2020-09-25
      • 1970-01-01
      • 2012-02-19
      • 2013-02-19
      • 2010-11-24
      • 1970-01-01
      相关资源
      最近更新 更多