【问题标题】:How to create a dataframe from a string key=value delimited by ";"如何从由“;”分隔的字符串 key=value 创建数据框
【发布时间】:2019-01-17 22:54:38
【问题描述】:

我有一个结构如下的 Hive 表:

我需要读取字符串字段,分解键并变成 Hive 表列,最终表应该是这样的:

很重要,字符串中键的数量是动态的,键的名称也是动态的

尝试使用 Spark SQL 读取字符串,使用基于所有字符串的模式创建数据帧,并使用 saveAsTable() 函数将数据帧转换为 hive 最终表,但不知道如何执行此操作

有什么建议吗?

【问题讨论】:

  • 您在用分号拆分键和值方面取得了多大进展?
  • 目前使用 Hive 查询,但是速度很慢。

标签: apache-spark hive apache-spark-sql


【解决方案1】:

一个天真的(假设唯一的(code, date) 组合并且在string 中没有嵌入=;)可能看起来像这样:

import org.apache.spark.sql.functions.{explode, split}

val df = Seq(
    (1, 1, "key1=value11;key2=value12;key3=value13;key4=value14"),
    (1, 2, "key1=value21;key2=value22;key3=value23;key4=value24"),
    (2, 4, "key3=value33;key4=value34;key5=value35")
).toDF("code", "date", "string")

val bits = split($"string", ";")
val kv = split($"pair", "=")

df
  .withColumn("bits", bits)  // Split column by `;`
  .withColumn("pair", explode($"bits"))  // Explode into multiple rows
  .withColumn("key", kv(0))  // Extract key
  .withColumn("val", kv(1))  // Extract value 
  // Pivot to wide format
  .groupBy("code", "date")
  .pivot("key")
  .agg(first("val"))

// +----+----+-------+-------+-------+-------+-------+
// |code|date|   key1|   key2|   key3|   key4|   key5|
// +----+----+-------+-------+-------+-------+-------+
// |   1|   2|value21|value22|value23|value24|   null|
// |   1|   1|value11|value12|value13|value14|   null|
// |   2|   4|   null|   null|value33|value34|value35|
// +----+----+-------+-------+-------+-------+-------+

(code, date) 不是唯一时,可以轻松调整以处理这种情况,您可以使用UDF 处理更复杂的string 模式。

根据您使用的语言和列数,使用RDDDataset 可能会更好。还值得考虑放弃完整的 explode / pivot 以支持 UDF。

val parse = udf((text: String) => text.split(";").map(_.split("=")).collect {
  case Array(k, v) => (k, v)
}.toMap)

val keys = udf((pairs: Map[String, String]) => pairs.keys.toList)

// Parse strings to Map[String, String]
val withKVs = df.withColumn("kvs", parse($"string"))

val keys = withKVs
  .select(explode(keys($"kvs"))).distinct // Get unique keys
  .as[String] 
  .collect.sorted.toList // Collect and sort

// Build a list of expressions for subsequent select
val exprs = keys.map(key => $"kvs".getItem(key).alias(key)) 

withKVs.select($"code" :: $"date" :: exprs: _*)

在 Spark 1.5 中你可以尝试:

val keys = withKVs.select($"kvs").rdd
  .flatMap(_.getAs[Map[String, String]]("kvs").keys)
  .distinct
  .collect.sorted.toList

【讨论】:

  • 可惜我用的是spark 1.5,pivot功能不可用,不然用不上?
  • 用RDD替换Dataset后的第二种方法应该可以正常工作。
  • .select 不适用于 RDD,我必须在哪里更改?
  • withKVs.select($"kvs").rdd.flatMap(_.getAs[Map[String, String]]("kvs").keys).distinct - 收集,其余的和以前一样。
  • 抱歉,我从来没有使用过 scala 我停止了在线:// Parse strings to Map [String, String] val = withKVs df.withColumn ("kvs" parse ($ "string")) 接下来会发生什么?
猜你喜欢
  • 2011-06-13
  • 1970-01-01
  • 2018-09-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-03-25
相关资源
最近更新 更多