【问题标题】:Spark-Scala Question on Using Case Class DemiliterSpark-Scala 关于使用 Case Class Demiliter 的问题
【发布时间】:2021-08-18 12:53:08
【问题描述】:

我有一个名为stockmarket.txt的文本文件:

2012-01-03,59.970001,61.060001,59.869999,60.330002,12668800,52.619234999999996,2012-01-04,60.209998999999996,60.349998,59.470001,59.709998999999996,9593300,52.078475,2012-01-05,59.349998,59.619999,58.369999,59.419998,12768200,51.825539,2012-01-06,59.419998,59.450001,58.869999,59.0,8069400,51.45922,2012-01-09,59.029999,59.549999,58.919998,59.18,6679300,51.616215000000004,2012-01-10,59.43,59.709998999999996,58.98,59.040001000000004,6907300,51.494109,2012-01-11,59.060001,59.529999,59.040001000000004,59.400002,6365600,51.808098,2012-01-12,59.790001000000004,60.0,59.400002,59.5,7236400,51.895315999999994

我想在每出现第 7 个逗号作为分隔符后将该行分隔成一个新行并显示结果。将有 7 列命名为 Date,Open,High,Low,Close,Volume,AdjClose 应使用数据进行标记。必要时使用案例类。

【问题讨论】:

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


    【解决方案1】:
    
    case class Stock(date: String, open: String, high: Double, low: Double, close: Double, volume: Long, adjClose: Double)
    
    ds.withColumn("value", 
        regexp_replace('value, "([^,]*),([^,]*),([^,]*),([^,]*),([^,]*),([^,]*),([^,]*),", "$1,$2,$3,$4,$5,$6,$7\n"))
      .withColumn("value", explode(split('value, "\n")))
      .withColumn("value", split('value, ",")).as[Seq[String]]
      .map(a => Stock(a.head, a(1), a(2).toDouble, a(3).toDouble, a(4).toDouble, a(5).toLong, a(6).toDouble))
    

    产生

    +----------+------------------+------------------+------------------+------------------+--------+------------------+
    |date      |open              |high              |low               |close             |volume  |adjClose          |
    +----------+------------------+------------------+------------------+------------------+--------+------------------+
    |2012-01-03|59.970001         |61.060001         |59.869999         |60.330002         |12668800|52.619234999999996|
    |2012-01-04|60.209998999999996|60.349998         |59.470001         |59.709998999999996|9593300 |52.078475         |
    |2012-01-05|59.349998         |59.619999         |58.369999         |59.419998         |12768200|51.825539         |
    |2012-01-06|59.419998         |59.450001         |58.869999         |59.0              |8069400 |51.45922          |
    |2012-01-09|59.029999         |59.549999         |58.919998         |59.18             |6679300 |51.616215000000004|
    |2012-01-10|59.43             |59.709998999999996|58.98             |59.040001000000004|6907300 |51.494109         |
    |2012-01-11|59.060001         |59.529999         |59.040001000000004|59.400002         |6365600 |51.808098         |
    |2012-01-12|59.790001000000004|60.0              |59.400002         |59.5              |7236400 |51.895315999999994|
    +----------+------------------+------------------+------------------+------------------+--------+------------------+
    

    为了测试你的字符串,我做了:

    val str = """<your sample string in here>"""
    val ds = Seq(str).toDS  // you'd use `spark.read.text("stockmarket.txt")` or similar
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-11-25
      • 1970-01-01
      • 2021-02-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-11-25
      • 2016-02-16
      相关资源
      最近更新 更多