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
相关产品推荐
相关产品推荐

