如何查询Driver/Worker节点配置、识别UDF运行节点并均衡分配PySpark UDF
问题1:获取Worker/Driver节点配置,识别UDF是否在Driver运行
获取节点配置
Driver节点配置
- 通过Spark配置获取核心、内存等参数:
# 获取Driver核心数 print(spark.conf.get("spark.driver.cores")) # 获取Driver内存 print(spark.conf.get("spark.driver.memory")) - 在Databricks Notebook中执行Shell命令查看硬件细节:
# 查看CPU信息 cat /proc/cpuinfo # 查看内存信息 free -m - 也可以通过Spark UI的Environment页面,查看所有Driver相关的配置项。
Worker节点配置
- 通过Spark配置获取executor(默认每个Worker对应一个executor)的参数:
# Worker节点数(executor实例数) print(spark.conf.get("spark.executor.instances")) # 每个Worker的核心数 print(spark.conf.get("spark.executor.cores")) # 每个Worker的内存 print(spark.conf.get("spark.executor.memory")) - 查看Spark UI的Executors页面,可看到每个Worker的IP、资源使用、运行状态等详细信息。
识别UDF是否在Driver运行
UDF默认在Worker的executor中执行,但小数据集可能触发本地执行(Driver上运行),可通过以下方式验证:
- 在UDF中嵌入节点信息采集逻辑:
import socket from pyspark.sql.functions import udf from pyspark.sql.types import StringType def get_host_info(): # 返回主机名和IP return f"{socket.gethostname()} | {socket.gethostbyname(socket.gethostname())}" host_info_udf = udf(get_host_info, StringType()) - 将UDF应用到目标DataFrame,查看结果中的节点信息:
如果结果中出现Driver的主机名/IP,说明该UDF在Driver上运行。df_with_host = df.withColumn("host_info", host_info_udf()) df_with_host.show(truncate=False) - 查看Spark UI的Jobs→Stages页面,查看Task的运行节点IP,若对应Driver节点,则UDF在Driver执行。
问题2:确保n个UDF均匀分配到n个Worker,且同时处理
实现均匀分配的步骤
调整DataFrame分区数匹配Worker数
Spark的Task按分区分配,要让每个Worker处理一个UDF,需创建分区数等于Worker数(n)的DataFrame:# 创建n行数据,重分区为n个分区 df = spark.range(n).repartition(n)配置集群资源保证每个Worker只跑一个Task
确保每个Worker仅运行一个executor(默认配置),且executor占满Worker资源:- 在集群配置中设置
spark.executor.instances = n(Worker节点数) - 设置
spark.executor.cores为Worker的总核心数,避免单个Worker同时运行多个Task
- 在集群配置中设置
在UDF中嵌入gRPC调用逻辑
注意gRPC客户端需在UDF内部初始化,避免Driver端的实例无法序列化到Worker:import grpc import socket from pyspark.sql.functions import udf from pyspark.sql.types import StringType def grpc_call_udf(input_id): # 在UDF内部初始化gRPC客户端 channel = grpc.insecure_channel("external-api:50051") # 执行gRPC调用逻辑(此处替换为实际调用代码) response = "success" return f"{response} from {socket.gethostname()} for id {input_id}" grpc_udf = udf(grpc_call_udf, StringType())然后将UDF应用到分区后的DataFrame:
result_df = df.withColumn("grpc_result", grpc_udf(df.id)) result_df.show()
验证Worker同时处理UDF的方式
- 查看节点分布结果
使用问题1中的host_info_udf,将其与gRPC UDF一同应用,查看host_info列的结果,每个Worker的主机名应仅出现一次,说明每个Worker分配到一个UDF。 - 实时监控Spark UI
- 打开Spark UI的Executors页面,在任务运行时,观察每个Executor的Active Tasks数是否为1,且所有Executor均有Active Task。
- 查看Jobs→Stages页面,查看Task的运行节点分布,确认每个Worker对应一个Task。
- 查看Executor日志
在Databricks的Job Run Details页面,进入每个Executor的日志,搜索gRPC调用的日志信息,确认同一时间点每个Worker都有执行记录。
内容的提问来源于stack exchange,提问作者boring-coder
相关产品推荐
相关产品推荐

