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

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â

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:58:36