【发布时间】:2022-01-11 01:08:25
【问题描述】:
我正在尝试转换通过火花流获得的输入,以便从中创建数据帧。基本上我会收到一个我想要从中提取数据的 json 字符串列表。
注意:我将 json 字符串简化为仅适用于一般概念的 coords 对象。
我得到的输入:
["{\"coord\":{\"lon\":10.0217,\"lat\":53.5281}}", "{\"coord\":{"lon\":10.1169,\"纬度\":53.6522}}", "{\"坐标\":...."]
我要创建的数据框以将其保存到数据库中:
+----------+----------+
|lon |lat |
+----------+----------+
| 10.0217| 53.5281|
| 10.1169| 53.6522|
| ... | ... |
+----------+----------+
到目前为止,我设法用字符串数组替换了引号。 我试图展平数组:
result = df \
.selectExpr("Cast(value AS STRING) as json") \
.withColumn("json", f.regexp_replace('json', '\\\\"', '"')) \
.withColumn("json", f.flatten(f.col("json"))) \
.select("json")
错误:
pyspark.sql.utils.AnalysisException: 无法解析 'flatten(
json)' 由于数据类型不匹配:参数应该是数组数组, 但是'json'是字符串类型的。;;
然后我尝试使用json.loads 加载数组,但我无法从 Spark 流式传输中调用此函数。
那么如何从这个输入中提取数据呢?
【问题讨论】:
-
lon前面好像少了一个反斜杠,对吗? -
@BrendanA 已更正。我在为 stackoverflow 双重转义时错过了它。谢谢指出!
标签: apache-spark pyspark spark-streaming