【问题标题】:Pyspark Structred Streaming Parse Nested JsonPyspark结构化流解析嵌套Json
【发布时间】:2020-04-29 03:19:08
【问题描述】:

我的项目是,将 json 写入 Kafka 主题并从 kafka 主题中读取 json,最后下沉一个 csv。一切正常。但是一些关键是嵌套的 json。如何解析 json 中的列表?

示例 Json:

{"a": "test", "b": "1234", "c": "temp", "d": [{"test1": "car", "test2": 345}, {"test3": "animal", "test4": 1}], "e": 50000}

你可以在下面看到我的代码。

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.types import *
import pyspark.sql.functions as func
spark = SparkSession.builder\
                    .config('spark.jars.packages', 'org.apache.spark:spark-sql-kafka-0-10_2.11:2.3.0') \
                    .appName('kafka_stream_test')\
                    .getOrCreate()
ordersSchema = StructType() \
        .add("a", StringType()) \
        .add("b", StringType()) \
        .add("c", StringType()) \
        .add("d", StringType())\
        .add("e", StringType())

df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "test") \
    .load()\


df_query = df \
    .selectExpr("cast(value as string)") \
    .select(func.from_json(func.col("value").cast("string"),ordersSchema).alias("parsed"))\
    .select("parsed.a","parsed.b","parsed.c","parsed.d","parsed.e","parsed.f")\

df_s = df_query \
    .writeStream \
    .format("console") \
    .outputMode("append") \
    .trigger(processingTime = "1 seconds")\
    .start()


aa = df_query \
    .writeStream \
    .format("csv")\
    .trigger(processingTime = "5 seconds")\
    .option("path", "/var/kafka_stream_test_out/")\
    .option("checkpointLocation", "/var/kafka_stream_test_out/chk") \
    .start()


df.printSchema()
df_s.awaitTermination()
aa.awaitTermination()

谢谢!

【问题讨论】:

    标签: python json apache-spark pyspark spark-structured-streaming


    【解决方案1】:

    “d”列的架构错误。它需要是一个 ArrayType。请查看等效的 Scala 代码,您可以将其转换为 Python。

        val schema = new StructType().add("a",StringType)
          .add("b",StringType)
          .add("c",StringType)
          .add("d",ArrayType(new StructType().add("test1",StringType).add("test2",StringType)))
          .add("e",StringType)
    

    json在“d”的每一行都有不同的列名。我假设这是一个错字,字段是“test1”和“test2”

    【讨论】:

    • 谢谢。我解决了,但我还有另一个问题。问题是,CSV 数据源不支持 array 数据类型。
    猜你喜欢
    • 1970-01-01
    • 2018-12-01
    • 2018-03-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-28
    • 1970-01-01
    相关资源
    最近更新 更多