【问题标题】:How to define schema for JSON records with timestamp (from Kafka) using (Py)Spark Structured Streaming? - null values shown如何使用(Py)Spark Structured Streaming为带有时间戳(来自Kafka)的JSON记录定义模式? - 显示空值
【发布时间】:2020-03-08 17:02:02
【问题描述】:

问题是我在使用 PySpark 阅读 Kafka 消息后得到了 null 值。

我使用 Spark 2.3.1 / Scala 2.11.12

我的代码:

allData = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", "localhost:9092") \
  .option("subscribe", "mysql.login") \
  .option("startingOffsets", "earliest") \
  .load()

df = allData.selectExpr("cast(value as string)", "timestamp", "topic" )

detailSchema = StructType() \
    .add("username", StringType()) \
    .add("login_time", DateType())

df2 = df.select(from_json(col('value'), detailSchema).alias('data'), 'timestamp', 'topic')

writeStream3 = df2 \
    .writeStream \
    .trigger(processingTime= '4 seconds') \
    .format('console') \
    .outputMode('update') \
    .start()

writeStream3.awaitTermination()    

kafka-console-consumer.sh阅读的消息如下:

$ kafka-console-consumer.sh \
    --bootstrap-server 127.0.0.1:9092 \
    --topic mysql.login \
    --from-beginning
{"username":"hello kitty","login_time":1572866627000}
{"username":"chitara","login_time":1572867234000}
{"username":"hello kitty","login_time":1572868094000}

但是,当我尝试阅读消息时,我看不到值。它在以下行之后显示为null:

df2 = df.select(from_json(col('value'), detailSchema).alias('data'), 'timestamp', 'topic')

我的代码的输出是:

+--------------------+--------------------+-----------+
|               value|           timestamp|      topic|
+--------------------+--------------------+-----------+
|{"username":"hell...|2019-11-12 13:55:...|mysql.login|
|{"username":"chit...|2019-11-12 13:55:...|mysql.login|
|{"username":"hell...|2019-11-12 13:55:...|mysql.login|
|{"username":"leon...|2019-11-12 13:55:...|mysql.login|
|{"username":"chit...|2019-11-12 13:55:...|mysql.login|
...

+----+--------------------+-----------+
|data|           timestamp|      topic|
+----+--------------------+-----------+
|null|2019-11-12 13:55:...|mysql.login|
|null|2019-11-12 13:55:...|mysql.login|
|null|2019-11-12 13:55:...|mysql.login|
|null|2019-11-12 13:55:...|mysql.login|
|null|2019-11-12 13:55:...|mysql.login|
...

+--------+-----+
|username|count|
+--------+-----+
|    null|  242|
+--------+-----+

我认为这个问题与解析有关,这就是为什么我在 from_json 函数之后看到 null 值的原因。为什么?如何解决?

【问题讨论】:

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


    【解决方案1】:

    tl;dr 将TimestampType 用于login_time。


    由于login_time 是时间戳,您应该使用正确的类型,例如TimestampType 或 LongType。

    来自official documentation:

    在不可解析字符串的情况下返回null。

    这正是您从 from_json 得到的,因为架构与输入行不匹配。

    【讨论】:

      猜你喜欢
      • 2017-10-24
      • 2023-03-08
      • 1970-01-01
      • 2018-07-23
      • 2021-04-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多