如何监控Databricks中PySpark applyInPandas()的计算进度?
groupBy().applyInPandas()计算进度的可行方法 Spark UI Stage监控
直接打开Databricks集群的Spark UI,查看Stage详情页。applyInPandas会把每个分组的计算拆成独立Task,这里能看到总Task数、已完成Task数,还能点击单个Task查看运行时长、状态。通过已完成数和总数的比例,就能直观掌握进度。在Python函数中加日志输出
修改你的python_function,每次处理完一个分组就打印日志,示例代码:import datetime def python_function(df): group_key = df["variable1"].iloc[0] # 你的核心计算逻辑 print(f"分组 {group_key} 处理完成,时间:{datetime.datetime.now()}") return result_df之后在Databricks的作业日志或集群Driver日志里,就能实时看到已完成的分组,间接判断进度。要是需要更规范的日志,也可以用Python的
logging模块输出到Databricks日志系统。自定义进度追踪表
在Databricks Metastore里建一张进度记录表:CREATE TABLE IF NOT EXISTS group_processing_progress ( group_key STRING, status STRING, start_time TIMESTAMP, end_time TIMESTAMP )然后在
python_function里嵌入状态更新逻辑:from pyspark.sql import SparkSession def python_function(df): spark = SparkSession.getActiveSession() group_key = str(df["variable1"].iloc[0]) # 标记分组开始运行 spark.sql(f""" INSERT INTO group_processing_progress VALUES ('{group_key}', 'running', CURRENT_TIMESTAMP(), NULL) ON DUPLICATE KEY UPDATE status='running', start_time=CURRENT_TIMESTAMP() """) # 你的计算代码 # 标记分组完成 spark.sql(f""" UPDATE group_processing_progress SET status='completed', end_time=CURRENT_TIMESTAMP() WHERE group_key='{group_key}' """) return result_df之后随时查询这张表,比如用
SELECT COUNT(*) FROM group_processing_progress WHERE status='completed',对比总分组数就能算出进度,还能查看具体哪些分组还在运行。注意如果variable1是数字类型,要调整表中group_key的字段类型。Databricks作业任务监控(作业模式下)
如果是通过Databricks Jobs提交的任务,直接在作业页面查看Run的实时状态,这里会显示已完成Task数、总Task数,还能查看每个Task的日志。要是开启了作业监控指标,在Metrics页面还能看到CPU、内存的实时使用情况,辅助判断计算推进状态。
内容的提问来源于stack exchange,提问作者Nairolf

