Databricks中Spark Streaming读取Kafka金融数据写流无结果求助
解决Spark Streaming读取Kafka数据无输出问题
核心问题分析
你当前的代码仅读取了Kafka的原始元数据结构(包含key、value、topic、partition等字段),并未解析消息中存储的业务JSON数据;同时流查询启动后未等待执行,或消费起始位置配置不合理,导致无输出结果。
解决方案步骤
1. 解析Kafka消息的业务JSON数据
Kafka消息的value字段为二进制格式,需先转换为字符串,再解析为结构化DataFrame:
from pyspark.sql.functions import from_json, col from pyspark.sql.types import StructType, StructField, StringType # 定义与Alpha Vantage消息匹配的Schema schema = StructType([ StructField("From_Currency Code", StringType()), StructField("From_Currency Name", StringType()), StructField("To_Currency Code", StringType()), StructField("To_Currency Name", StringType()), StructField("Exchange Rate", StringType()), StructField("Last Refreshed", StringType()), StructField("Time Zone", StringType()), StructField("Bid Price", StringType()), StructField("Ask Price", StringType()) ]) # 读取Kafka并解析业务数据 df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", 'localhost:9092') \ .option("subscribe", topic) \ # 读取主题历史消息,若仅需新消息可改为"latest" .option("startingOffsets", "earliest") \ .load() \ # 将二进制value转为字符串 .select(col("value").cast(StringType()).alias("json_str")) \ # 解析JSON为结构化数据 .select(from_json(col("json_str"), schema).alias("data")) \ .select("data.*")
2. 正确启动并等待流查询执行
在Databricks中,流查询启动后需等待其处理数据,否则会立即终止:
query = df.writeStream \ .outputMode("append") \ .format("console") \ .queryName("fx") \ # 设置触发间隔,按需调整 .trigger(processingTime="5 seconds") \ .start() # 阻塞等待查询执行,直到手动停止或出错 query.awaitTermination()
3. 额外排查点
- 确认
localhost:9092在Databricks环境中可访问,若为远程Kafka需配置正确地址与网络规则。 - 检查
topic变量是否指向目标Kafka主题。 - 若使用
startingOffsets="latest",需确保查询启动后有新消息写入Kafka,否则无输出。
内容的提问来源于stack exchange,提问作者Hayel Douaâ
相关产品推荐
相关产品推荐

