PySpark与Databricks中addFile/SparkFiles.get触发文件未找到异常
核心原因是JDBC URL中的证书路径是Driver本地路径,而Executor执行Task时使用该URL,导致在Executor上找不到对应文件。你在Driver端通过SparkFiles.get()获取的路径仅在Driver节点有效,Executor节点的文件路径与Driver端不同,直接硬编码到JDBC URL中会导致Executor找不到文件。
错误信息中端口显示53,101属于错误日志的格式化问题(千位分隔符),实际端口53101是正确的,无需额外处理。
先确认证书文件是否已正确分发到所有Executor节点:
from pyspark import SparkFiles import os def check_cert_file(filename): file_path = SparkFiles.get(filename) return f"节点路径: {file_path}, 文件存在: {os.path.exists(file_path)}" # 用并行任务验证所有Executor节点 check_results = sc.parallelize(range(2)).map(lambda x: check_cert_file("dwt01_db2_ssl.arm")).collect() print("Executor证书文件检查结果:", check_results)
如果输出显示所有节点都存在文件,说明分发正常,问题出在路径传递;如果不存在,检查Driver端下载的文件是否在当前工作目录(s3_client.download_file的Filename是相对路径,需确保Driver能找到)。
方案1:自定义分区读取(推荐)
通过mapPartitions在每个Executor分区中动态获取证书路径,构建正确的JDBC连接:
from pyspark import SparkFiles import pandas as pd import jaydebeapi def read_db2_data(partition): # 在Executor节点动态获取证书路径 cert_path = SparkFiles.get("dwt01_db2_ssl.arm") # 构建Executor本地可用的JDBC URL jdbc_url = f"jdbc:db2://{hostname}:{port}/{database}:sslConnection=true;sslCertLocation={cert_path};" # 建立DB2连接并读取数据 conn = jaydebeapi.connect( driver_name="com.ibm.db2.jcc.DB2Driver", url=jdbc_url, credentials=[user_name, password] # 若驱动jar未在classpath中,需添加参数: jars="/path/to/db2jcc4.jar" ) df = pd.read_sql(query, conn) conn.close() # 返回数据行 return df.itertuples(index=False, name=None) # 创建空RDD触发分区执行,可根据需求设置分区数 empty_rdd = sc.parallelize([], numSlices=4) result_df = empty_rdd.mapPartitions(read_db2_data).toDF()
此方法确保每个Executor使用本地的证书路径,彻底解决路径不一致问题。
方案2:预分发证书到集群节点
如果集群是固定的,可通过Ansible、Shell脚本等方式将证书文件同步到所有Executor节点的同一个绝对路径(如/opt/db2_certs/dwt01_db2_ssl.arm),然后在JDBC URL中直接使用该绝对路径:
jdbc_url = f"jdbc:db2://{hostname}:{port}/{database}:sslConnection=true;sslCertLocation=/opt/db2_certs/dwt01_db2_ssl.arm;"
此方法适合长期固定的证书配置,无需依赖Spark的文件分发机制。
方案3:升级DB2 JDBC驱动版本
你当前使用的驱动版本是4.27.25,尝试升级到最新版本(如4.31.14),新版本可能优化了文件路径的解析逻辑,解决跨节点路径问题。
内容的提问来源于stack exchange,提问作者eemilk

