You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何查询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上运行),可通过以下方式验证:

  1. 在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())
    
  2. 将UDF应用到目标DataFrame,查看结果中的节点信息:
    df_with_host = df.withColumn("host_info", host_info_udf())
    df_with_host.show(truncate=False)
    
    如果结果中出现Driver的主机名/IP,说明该UDF在Driver上运行。
  3. 查看Spark UI的Jobs→Stages页面,查看Task的运行节点IP,若对应Driver节点,则UDF在Driver执行。

问题2:确保n个UDF均匀分配到n个Worker,且同时处理

实现均匀分配的步骤

  1. 调整DataFrame分区数匹配Worker数
    Spark的Task按分区分配,要让每个Worker处理一个UDF,需创建分区数等于Worker数(n)的DataFrame:

    # 创建n行数据,重分区为n个分区
    df = spark.range(n).repartition(n)
    
  2. 配置集群资源保证每个Worker只跑一个Task
    确保每个Worker仅运行一个executor(默认配置),且executor占满Worker资源:

    • 在集群配置中设置spark.executor.instances = n(Worker节点数)
    • 设置spark.executor.cores为Worker的总核心数,避免单个Worker同时运行多个Task
  3. 在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. 查看节点分布结果
    使用问题1中的host_info_udf,将其与gRPC UDF一同应用,查看host_info列的结果,每个Worker的主机名应仅出现一次,说明每个Worker分配到一个UDF。
  2. 实时监控Spark UI
    • 打开Spark UI的Executors页面,在任务运行时,观察每个Executor的Active Tasks数是否为1,且所有Executor均有Active Task。
    • 查看Jobs→Stages页面,查看Task的运行节点分布,确认每个Worker对应一个Task。
  3. 查看Executor日志
    在Databricks的Job Run Details页面,进入每个Executor的日志,搜索gRPC调用的日志信息,确认同一时间点每个Worker都有执行记录。

内容的提问来源于stack exchange,提问作者boring-coder

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 00:38:16