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

如何查看并输出Spark Streaming的错误日志?

Spark Streaming 错误日志查看及报错排查指引

错误日志查看方法

可以查看完整的Spark Streaming错误日志,具体操作如下:

  • 检索Driver端全量日志:你当前贴出的仅为StreamingQueryException的表层摘要,完整的根因异常(包含Caused by层级的报错栈)会输出在Driver的stdout/stderr日志中,可直接搜索查询ID 6bc7f9c4-b8da-4cbf-8525-2ff720343ed2 定位完整报错信息。
  • 排查Executor端日志:如果异常是Executor侧抛出导致任务终止,可通过Spark UI的Jobs标签页找到对应失败的Job,进入对应Stage页面查看单Executor日志,也可直接到Executor部署节点的日志目录下检索关联任务ID的报错。

现有报错特征排查提示

你提供的报错信息显示:已提交的Kafka消费offset为950,对应分区最新可用offset已经达到6120479,说明任务在处理offset 950之后的批次数据时触发写入失败,优先排查以下两点:

  • 该offset区间的Kafka数据格式是否符合你的解析逻辑,是否存在脏数据导致解析异常
  • 数据写入目标的连通性、权限、资源配额是否正常

代码主动捕获异常方案

如果需要在代码中主动输出完整错误栈,可参考以下示例:

Scala 版本

import org.apache.spark.sql.streaming.StreamingQueryListener

spark.streams.addListener(new StreamingQueryListener {
  override def onQueryStarted(event: StreamingQueryListener.QueryStartedEvent): Unit = {}
  override def onQueryProgress(event: StreamingQueryListener.QueryProgressEvent): Unit = {}
  override def onQueryTerminated(event: StreamingQueryListener.QueryTerminatedEvent): Unit = {
    event.exception.foreach { ex =>
      logError(s"Streaming query ${event.id} 执行失败,完整异常栈: ", ex)
    }
  }
})

Python 版本

import traceback

query = df.writeStream.start()
try:
    query.awaitTermination()
except Exception as e:
    traceback.print_exc()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 11:30:01