You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark消费Kafka消息返回NULL值问题排查求助

问题

使用PySpark连接Kafka broker消费test-topic时,控制台输出的message字段始终为NULL,但用kafka-console-consumer.sh能看到正确的JSON消息。

消费者代码(consumer.py)

def consume_message():
    
    load_dotenv()

    spark = SparkSession.builder \
    .appName("KafkaReader") \
    .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0") \
    .getOrCreate()

    schema = StructType().add("message", StringType())

    df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9093") \
    .option("subscribe", "test-topic") \
    .option("startingOffsets", "latest") \
    .load()
    
    messages = df.select(from_json(col("value").cast("string"), schema).alias("data")) \
    .select("data.message")
    
    query = messages \
    .writeStream \
    .outputMode("append") \
    .format("console") \
    .start()
    query.awaitTermination()


if __name__ == '__main__':
    time.sleep(10)
    while True:        
        consume_message()
        time.sleep(60)

PySpark消费输出

2025-03-05 23:15:16 -------------------------------------------
2025-03-05 23:15:16 Batch: 25
2025-03-05 23:15:16 -------------------------------------------
2025-03-05 23:15:16 +-------+
2025-03-05 23:15:16 |message|
2025-03-05 23:15:16 +-------+
2025-03-05 23:15:16 |   NULL|
2025-03-05 23:15:16 +-------+
2025-03-05 23:15:16 
2025-03-05 23:16:18 -------------------------------------------
2025-03-05 23:16:18 Batch: 26
2025-03-05 23:16:18 -------------------------------------------
2025-03-05 23:16:18 +-------+
2025-03-05 23:16:18 |message|
2025-03-05 23:16:18 +-------+
2025-03-05 23:16:18 |   NULL|
2025-03-05 23:16:18 +-------+

Kafka控制台消费结果

通过以下命令直接消费能看到正确的JSON:

docker exec -it <container_id>  bash

cd opt/bitnami/kafka/bin

kafka-console-consumer.sh --bootstrap-server kafka:9093 --topic test-topic --from-beginning

输出:

{"date": "2025-03-05 22:03", "BTCUSDT": "90244.50000000", "ETHUSDT": "2230.69000000", "BNBUSDT": "595.95000000", "SOLUSDT": "145.04000000", "ADAUSDT": "0.98140000", "XRPUSDT": "2.49800000", "SHIBUSDT": "0.00001328", "DOGEUSDT": "0.20369000", "TONUSDT": "3.02600000"}
{"date": "2025-03-05 22:04", "BTCUSDT": "90233.10000000", "ETHUSDT": "2230.99000000", "BNBUSDT": "595.91000000", "SOLUSDT": "145.02000000", "ADAUSDT": "0.98330000", "XRPUSDT": "2.49890000", "SHIBUSDT": "0.00001330", "DOGEUSDT": "0.20379000", "TONUSDT": "3.02600000"}

已尝试的修改

修改schema定义为以下形式,结果仍为NULL:

schema = StructType([StructField("message", StringType())])

解决方案

问题根源

定义的schema期望Kafka消息的value是一个包含message字段的JSON结构(比如{"message": "实际内容"}),但实际Kafka中的消息本身就是顶层JSON,直接包含date、BTCUSDT等字段,没有message外层字段。因此from_json解析不匹配,返回NULL。

两种解决方式

方式1:直接获取原始JSON字符串

如果不需要将JSON解析为结构化数据,直接把value转成字符串即可:

messages = df.select(col("value").cast("string").alias("message"))

方式2:定义匹配实际消息的Schema

如果需要解析成结构化DataFrame,要根据实际消息的字段定义正确的Schema:

from pyspark.sql.types import StructType, StructField, StringType

# 匹配实际消息的字段结构
schema = StructType([
    StructField("date", StringType()),
    StructField("BTCUSDT", StringType()),
    StructField("ETHUSDT", StringType()),
    StructField("BNBUSDT", StringType()),
    StructField("SOLUSDT", StringType()),
    StructField("ADAUSDT", StringType()),
    StructField("XRPUSDT", StringType()),
    StructField("SHIBUSDT", StringType()),
    StructField("DOGEUSDT", StringType()),
    StructField("TONUSDT", StringType())
])

messages = df.select(from_json(col("value").cast("string"), schema).alias("data")) \
             .select("data.*")  # 展开所有字段

修改后即可正确解析Kafka中的消息,不再返回NULL。

内容的提问来源于stack exchange,提问作者KurczakChrupiacy2

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 23:28:20