如何获取Spark中单个转换操作(如drop、rename)的完成时间
如何获取Spark单个转换操作的执行时间
Spark里的drop、withColumnRenamed这类转换操作是惰性求值的——它们不会立即执行,只会构建计算逻辑的DAG(有向无环图),直到遇到show()、write()、count()这类行动操作时,才会把所有累积的转换一起提交执行。所以直接单独测单个转换的“完成时间”是不现实的,得通过触发轻量行动来间接测量每个转换对应的执行耗时。
方法一:手动计时+轻量行动触发
可以借助Python的time模块,在每个转换后执行一个轻量的行动(比如count(),不会产生大量输出),记录前后的时间差,以此近似该转换(结合之前未执行的转换)的耗时。
修正你代码里的中文引号问题后,示例代码如下:
import pyspark import time from pyspark.sql import SparkSession # 创建SparkSession data = [("Banana",1000,"USA","1"), ("Carrots",1500,"USA","2"), ("Beans",1600,"USA","3"), ("Orange",2000,"USA","4"),("Banana",400,"China","5"), ("Carrots",1200,"China","1"),("Beans",1500,"China","2"),("Orange",4000,"China","3"), ("Banana",2000,"Canada","4"),("Carrots",2000,"Canada","5"),("Beans",2000,"Mexico","6"),("Orange",2000,"USA","7")] columns= ["Product","Amount","Country","Id"] spark = SparkSession.builder.master("local[*]").getOrCreate() df = spark.createDataFrame(data = data, schema = columns) # 测量drop操作的耗时 start_time = time.time() df_drop = df.drop("Id") # 触发行动执行转换 df_drop.count() drop_time = time.time() - start_time print(f"drop操作执行耗时: {drop_time:.4f}秒") # 测量withColumnRenamed操作的耗时 start_time = time.time() df_rename = df_drop.withColumnRenamed("Product","Veggies") df_rename.count() rename_time = time.time() - start_time print(f"withColumnRenamed操作执行耗时: {rename_time:.4f}秒") # 后续操作 df_rename.write.csv("Output.csv") df_rename.show(truncate=False)
方法二:通过Spark UI查看阶段耗时
Spark运行时会默认启动Web UI(默认端口4040),你可以在应用运行期间访问http://localhost:4040,进入Stages页面查看每个执行阶段的详细耗时:
- 每个阶段对应DAG中的一组转换操作,你可以根据阶段的描述(比如包含
drop、rename相关的算子)匹配到对应的转换,查看其执行时间。 - 这种方法能看到更详细的执行细节,比如任务数、每个任务的耗时等。
注意事项
- 手动计时的结果会包含Spark调度任务的开销,不是纯转换的计算时间,但对于日常调试足够用。
- 如果多个转换连续执行且未插入行动操作,Spark会把它们合并到同一个阶段执行,此时无法单独区分每个转换的耗时,必须在每个转换后触发行动。
内容的提问来源于stack exchange,提问作者tejaswini teju
相关产品推荐
相关产品推荐

