如何查看并输出Spark Streaming的错误日志?
Spark Streaming 错误日志查看及报错排查指引
错误日志查看方法
可以查看完整的Spark Streaming错误日志,具体操作如下:
- 检索Driver端全量日志:你当前贴出的仅为
StreamingQueryException的表层摘要,完整的根因异常(包含Caused by层级的报错栈)会输出在Driver的stdout/stderr日志中,可直接搜索查询ID6bc7f9c4-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
相关产品推荐
相关产品推荐

