【问题标题】:Processing a list of json strings in Spark Streaming在 Spark Streaming 中处理 json 字符串列表
【发布时间】: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


【解决方案1】:

提供数组

arr = [
    "{\"coord\":{\"lon\":10.0217,\"lat\":53.5281}}", 
    "{\"coord\":{\"lon\":10.1169,\"lat\":53.6522}}", 
    ]

你可以用下面的代码得到想要的结果

from pyspark.sql import functions, types

df = (df.withColumn("lon", functions.regexp_extract("value", "(?<=lon\"\:)[0-9]+.[0-9]+", 0))
        .withColumn("lat", functions.regexp_extract("value", "(?<=lat\"\:)[0-9]+.[0-9]+", 0)))

df = df.select(df["lon"], df["lat"])

df.show()
+-------+-------+
|    lon|    lat|
+-------+-------+
|10.0217|53.5281|
|10.1169|53.6522|
+-------+-------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-04-08
    • 1970-01-01
    • 2023-02-26
    • 1970-01-01
    • 1970-01-01
    • 2016-01-10
    • 1970-01-01
    相关资源
    最近更新 更多