【问题标题】:Fixed length parsing in spark scalaspark scala中的固定长度解析
【发布时间】:2019-06-30 15:14:29
【问题描述】:

我已经创建了数据框,输入是这样的:

   +-----------------------------------+
   |value                              |
   +-----------------------------------+
   |1   PRE123                    21   |
   |2   TEST                      32   |
   |7   XYZ                       .7   |
   +-----------------------------------+

在下面的元数据信息的基础上,我们需要拆分上面的数据帧并创建一个新的数据帧,列名称为 id、name 和 class,它的开始和索引位置在这个 json 元数据中给出。

   {
    "columnName": "id",
    "start": 1,
    "end": 2
  },
  {
    "columnName": "name",
    "start": 5,
    "end": 10
  },
  {
    "columnName": "class",
    "start": 20,
    "end": 22
  }

输出:

  +---+------+-----+
  | id|  name|class|
  +---+------+-----+
  |  1|PRE123|   21|
  |  2|  TEST|   32|
  |  7|   XYZ|   .7|
  +---+------+-----+

为了加载 df,我创建了列表:

   list.+=(loadedDF.col("value").substr(fixedLength.getStart, (fixedLength.getEnd - fixedLength.getStart)).alias(fixedLength.getColumnName))

从这个列表中,我创建了数据框

var df: DataFrame = loadedDF.select(list: _*)

需要了解从元数据创建数据帧的顺序更好的方法。 由于创建的列表会将所有数据带到驱动程序节点。

【问题讨论】:

  • 嗨,Etisha,目前尚不清楚您要达到的目标。您能否提供一个简单的示例,其中包含输入和所需的输出?还有什么是固定长度/元数据,它们与您的要求有什么关系?
  • @AlexandrosBiratsis 请再看一遍帖子。
  • Etisha 您是否有严格的要求来保持所有数据的固定长度?下面的解决方案更加灵活,不需要提供任何固定的开始/结束变量。在这种情况下,您需要的唯一元数据是列名。然后代码将更易于维护和简单
  • 是的。我们严格要求为您的所有数据保持固定长度,并且它将是动态的。

标签: scala apache-spark parsing fixed-length-file


【解决方案1】:

如果我理解正确,您的要求是尝试从由任意数量的空格分隔的字符串中提取列。

这是一个带有 substr 函数的解决方案:

val df = Seq(
  ("1   PRE123         21"),
  ("2   TEST           32"),
  ("7   XYZ            .7"))
.toDF("value")

val colMetadata = Map("id" -> (1,2), "name" -> (5,10), "class" -> (20,22))

val columns = colMetadata.map { case (cname, meta) => 
  val len = meta._2 - meta._1   
  $"value".substr(meta._1, len).as(cname)
}.toSeq

df.select(columns:_*).show

当你没有可用的列边界时,一个通用的解决方案是使用 split 函数:

import org.apache.spark.sql.functions.split

val df = Seq(
  ("1   PRE123         21"),
  ("2   TEST           32"),
  ("7   XYZ            .7"))
.toDF("value")

val colNames = Seq("id", "name", "class")

val columns = colNames.zipWithIndex.map { case (cname, idx) =>
      split($"value", "\\s+").getItem(idx).as(cname)
}

df.select(columns:_*).show

输出:

+---+------+-----+
| id|  name|class|
+---+------+-----+
|  1|PRE123|   21|
|  2|  TEST|   32|
|  7|   XYZ|   .7|
+---+------+-----+

请注意,我使用\\s+ 作为分隔符。这表示一个或多个空格的正则​​表达式。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-30
    • 2011-06-22
    • 2012-05-17
    相关资源
    最近更新 更多