【问题标题】:Reading array fields in Spark 2.2在 Spark 2.2 中读取数组字段
【发布时间】:2018-06-20 02:24:33
【问题描述】:

假设您有一堆数据,其行如下所示:

{
    'key': [
        {'key1': 'value11', 'key2': 'value21'},
        {'key1': 'value12', 'key2': 'value22'}
    ]
}

我想将其读入 Spark Dataset。一种方法如下:

case class ObjOfLists(k1: List[String], k2: List[String])
case class Data(k: ObjOfLists)

那么你可以这样做:

sparkSession.read.json(pathToData).select(
    struct($"key.key1" as "k1", $"key.key2" as "k2") as "k"
)
.as[Data]

这很好用,但有点破坏数据;毕竟在数据中'key' 指向对象列表而不是列表对象。换句话说,我真正想要的是:

case class Obj(k1: String, k2: String)
case class DataOfList(k: List[Obj])

我的问题:我可以在select 中输入一些其他语法,从而允许将生成的Dataframe 转换为Dataset[DataOfList]


我尝试使用与上述相同的select 语法,结果:

线程“主”org.apache.spark.sql.AnalysisException 中的异常:需要一个数组字段但得到了struct<k1:array<string>,k2:array<string>>;

所以我也试过了:

sparkSession.read.json(pathToData).select(
    array(struct($"key.key1" as "k1", $"key.key2" as "k2")) as "k"
)
.as[DataOfList]

这个编译运行了,但是数据看起来像这样:

DataOfList(List(Obj(org.apache.spark.sql.catalyst.expressions.UnsafeArrayData@bb2a5516,org.apache.spark.sql.catalyst.expressions.UnsafeArrayData@bec5e4a7)))

还有其他想法吗?

【问题讨论】:

  • 当然,一种解决方法是读取数据,然后应用map 将其转换为正确的形式,但这可能会有点笨拙 - 特别是如果许多字段中只有一个导致一个问题。

标签: json apache-spark apache-spark-dataset


【解决方案1】:

只需重铸数据以反映预期的名称:

case class Obj(k1: String, k2: String)
case class DataOfList(k: Seq[Obj])

val text = Seq("""{
  "key": [
    {"key1": "value11", "key2": "value21"},
    {"key1": "value12", "key2": "value22"}
  ]
}""").toDS

val df = spark.read.json(text)

df
  .select($"key".cast("array<struct<k1:string,k2:string>>").as("k"))
  .as[DataOfList]
  .first
DataOfList(List(Obj(value11,value21), Obj(value12,value22)))

对于无关的对象,您可以在读取时定义架构:

val textExtended = Seq("""{
  "key": [
    {"key0": "value01", "key1": "value11", "key2": "value21"},
    {"key1": "value12", "key2": "value22", "key3": "value32"}
  ]
}""").toDS

val schemaSubset = StructType(Seq(StructField("key", ArrayType(StructType(Seq(
  StructField("key1", StringType),
  StructField("key2", StringType))))
)))

val df = spark.read.schema(schemaSubset).json(textExtended)

然后像以前一样继续。

【讨论】:

  • 这很有帮助 - 谢谢!但是,我不满意:从实验看来,对cast 的调用中的字符串必须完全指定数组中对象的架构,而普通的select 语法允许您只获取您想要的内容.即使在这些对象中有"key3""key4" 等(但您只关心"key1""key2"),您的解决方案是否有一个简单的修改?除了为 cast 创建详细架构之外?
  • 好的,这看起来满足了我的所有需求 - 我对解决方案的冗长程度感到惊讶,但回想起来,我想这种做法违背了@987654331 @API 有点。谢谢!
  • DataFrames 不谈,我认为这里有些冗长是不可避免的,因为我们需要在嵌套结构上映射结构和名称。例如,如果您使用泛型,则必须同时使用输入和输出结构。如果事情是平的,那会更容易df.select(explode($"key")).select($"col.key1" as "k1", $"col.key2" as "k2").as[Obj],如果名字匹配,更简单case class ObjRaw(key1: String, key2: String); df.select(explode($"key")).select($"col.*").as[ObjRaw]
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-05-27
  • 1970-01-01
  • 2023-03-14
  • 1970-01-01
  • 1970-01-01
  • 2020-11-20
相关资源
最近更新 更多