【发布时间】:2021-03-07 03:59:36
【问题描述】:
使用我的 Scala HTTP 客户端,我从 API GET 调用中检索到 JSON 格式的响应。
我的最终目标是将此JSON 内容写入AWS S3 存储桶,以便在运行简单的AWS Glue 爬虫时将其作为RedShift 上的表提供。
我的想法是解析这个JSON 消息并以某种方式转换为Spark DataFrame,以便稍后我可以将其以.csv、.parquet 或其他格式保存到我喜欢的S3 位置。
JSON 文件如下所示
{
"response": {
"status": "OK",
"start_element": 0,
"num_elements": 100,
"categories": [
{
"id": 1,
"name": "Airlines",
"is_sensitive": false,
"last_modified": "2010-03-19 17:48:36",
"requires_whitelist_on_external": false,
"requires_whitelist_on_managed": false,
"is_brand_eligible": true,
"requires_whitelist": false,
"whitelist": {
"geos": [],
"countries_and_brands": []
}
},
{
"id": 2,
"name": "Apparel",
"is_sensitive": false,
"last_modified": "2010-03-19 17:48:36",
"requires_whitelist_on_external": false,
"requires_whitelist_on_managed": false,
"is_brand_eligible": true,
"requires_whitelist": false,
"whitelist": {
"geos": [],
"countries_and_brands": []
}
}
],
"count": 148,
"dbg_info": {
"warnings": [],
"version": "1.18.1621",
"output_term": "categories"
}
}
}
我想映射到 Dataframe 的内容是 "categories" JSON Array 包含的内容。
我已经设法使用json4s.JsonMethods方法parse这样解析消息:
val parsedJson = parse(request) \\ "categories"
获得以下内容:
output: org.json4s.JValue = JArray(List(JObject(List((id,JInt(1)), (name,JString(Airlines)), (is_sensitive,JBool(false)), (last_modified,JString(2010-03-19 17:48:36)), (requires_whitelist_on_external,JBool(false)), (requires_whitelist_on_managed,JBool(false)), (is_brand_eligible,JBool(true)), (requires_whitelist,JBool(false)), (whitelist,JObject(List((geos,JArray(List())), (countries_and_brands,JArray(List()))))))), JObject(List((id,JInt(2)), (name,JString(Apparel)), (is_sensitive,JBool(false)), (last_modified,JString(2010-03-19 17:48:36)), (requires_whitelist_on_external,JBool(false)), (requires_whitelist_on_managed,JBool(false)), (is_brand_eligible,JBool(true)), (requires_whitelist,JBool(false)), (whitelist,JObject(List((geos,JArray(List())), (countries_and_brands,JArray(List()))))))))
但是,我完全不知道如何进行。我什至尝试过使用另一个名为 uJson 的 Scala 库:
val json = (ujson.read(request))
val tuples = json("response")("categories").arr /* <-- categories is an array */ .map { item =>
(item("id"), item("name"))
这次我只解析了两个字段进行测试,但这应该不会有太大变化。因此,我得到了以下结构:
tuples: scala.collection.mutable.ArrayBuffer[(ujson.Value, ujson.Value, ujson.Value, ujson.Value)] = ArrayBuffer((1,"Airlines",false,"2010-03-19 17:48:36"), (2,"Apparel",false,"2010-03-19 17:48:36"))
但是,这次我也不知道如何继续前进,我尝试的一切都会导致错误,主要与格式不兼容有关。
请随意提出任何其他方法来实现我的目标,即使它完全改变了我的工作流程。我宁愿好好学点东西。谢谢
【问题讨论】:
标签: json scala apache-spark apache-spark-sql