Apache Spark作业中是否可获取dataset计数并实现异常自动通知?
Spark Join行数爆炸问题自动化排查方案
1. 作业异常行数自动通知实现
- 基于
SparkListener自定义监听器,重写onStageCompleted方法,直接从StageInfo的outputRecords参数获取每个Stage的输出行数 - 提前按业务场景配置不同作业的合理行数阈值,关联作业的
spark.app.name或自定义标签匹配对应开发负责人,触发阈值时直接调用内部告警通道推送通知 - 可以通过全局配置
spark.extraListeners加载自定义监听器,无需业务代码修改,对开发透明
2. 无侵入式Dataset指标日志打印
- 用上述自定义监听器在Stage/Job结束时,自动将对应算子类型、输入输出行数、Shuffle数据量等指标统一打印到集群日志服务,业务侧无需手动添加
count()类统计代码 - 也可注册
QueryExecutionListener,在每一个Spark SQL/Dataset的action操作执行完成后,自动拉取查询统计指标写入日志,粒度可细化到单个Dataset操作维度
3. 执行计划自动落盘存储
- 同样通过
QueryExecutionListener的onSuccess回调,调用qe.executedPlan.toString()即可获取完整物理执行计划,可搭配作业ID、运行时长、异常指标等信息一并存入运维数据库或日志平台,后续排查可直接按作业维度检索 - 若需要携带实际运行metrics的执行计划,也可在作业结束后调用Spark REST API拉取对应Application的全量Stage信息和执行计划,统一归档存储
上述所有方案都可以通过集群全局配置生效,不需要业务开发修改现有代码,完全零侵入实现全链路自动化排查能力。
内容的提问来源于stack exchange,提问作者MiamiBeach
相关产品推荐
相关产品推荐

