【发布时间】:2019-03-18 16:39:28
【问题描述】:
根据面包屑,我拼凑了一些似乎可以做我想做的事情的代码:在后台运行,查看正在进行的作业,然后收集...任何可用的信息:
def do_background_monitoring(sc: pyspark.context.SparkContext):
thread = threading.Thread(target=monitor, args=[sc])
thread.start()
return thread
def monitor(sc: pyspark.context.SparkContext):
job_tracker: pyspark.status.StatusTracker = sc.statusTracker() # should this go inside the loop?
while True:
time.sleep(1)
for job_id in job_tracker.getActiveJobsIds():
job: pyspark.status.SparkJobInfo = job_tracker.getJobInfo(job_id)
stages = job.stageIds
# ???
但是,这就是我遇到死胡同的地方。根据the docs,stageIds 是int[],显然是 py4j 或其他不知道如何处理它的东西? (py4j claims otherwise...)
ipdb> stages
JavaObject id=o34
ipdb> stages.
equals notify wait
getClass notifyAll
hashCode toString
ipdb> stages.toString()
'[I@4b1009f3'
这是一个死胡同吗?还有其他方法可以实现这一目标吗?如果我愿意 愿意编写 scala 来执行此操作,我是否可以将这部分保存在 Scala 中并将其余部分保存在 Python 中?
【问题讨论】:
标签: python-3.x apache-spark pyspark