PySpark应用调用API时UserPrincipalNotFoundException的解决方法
问题:Databricks PySpark调用API后使用parallelize报错UserPrincipalNotFoundException
代码示例
def api_json(spark, api_param): try: token = generate_a_token() token_headers = {'Authorization': f"Bearer {token}"} response = requests.get(f'https://api_url/?api_end_point={api_param}', headers=token_headers) json_data = spark.sparkContext.parallelize([response.text]) df = spark.read.json(json_data) df.show() return df except Exception as error: traceback.print_exc() def generate_a_token(): token_data = { some_key1: some_value1, some_key2: some_value2 } headers = {"Content-type": "application/x-www-form-urlencoded"} data = bytes(urlencode(token_data).encode()) req = request.Request(url, data, headers) try: response = request.urlopen(req) json_response = json.load(response) token = json_response["access_token"] print(token) except Exception as e: print('Some Exception') return token
触发的异常
py4j.protocol.Py4JJavaError: An error occurred while calling o305.getTrustedPath.
: java.nio.file.attribute.UserPrincipalNotFoundException
完整异常栈
Traceback (most recent call last): File "/local_disk0/tmp/spark-some-path/pyspark_template-1.0.0-py3-none-any.whl/src/tests/Api_Response.py", line 23, in api_json json_data = spark.sparkContext.parallelize([response.text]) File "/databricks/spark/python/pyspark/context.py", line 562, in parallelize jrdd = session._create_rdd_from_local_trusted(c, numSlices, serializer) File "/databricks/spark/python/pyspark/sql/session.py", line 516, in _create_rdd_from_local_trusted temp_data_path = self._write_to_trusted_path(data, serializer) File "/databricks/spark/python/pyspark/sql/session.py", line 485, in _write_to_trusted_path temp_dir = self._jsparkSession.getTrustedPath() File "/databricks/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py", line 1305, in __call__ answer, self.gateway_client, self.target_id, self.name) File "/databricks/spark/python/pyspark/sql/utils.py", line 127, in deco return f(*a, **kw) File "/databricks/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py", line 328, in get_return_value format(target_id, ".", name), value) py4j.protocol.Py4JJavaError: An error occurred while calling o305.getTrustedPath. : java.nio.file.attribute.UserPrincipalNotFoundException at sun.nio.fs.UnixUserPrincipals.lookupName(UnixUserPrincipals.java:155) at sun.nio.fs.UnixUserPrincipals.lookupUser(UnixUserPrincipals.java:170) at sun.nio.fs.UnixFileSystem$LookupService$1.lookupPrincipalByName(UnixFileSystem.java:330) at org.apache.spark.sql.TrustedPathHelper.createTrustedTempDir(TrustedPathHelper.scala:86) at org.apache.spark.sql.TrustedPathHelper.createTrustedTempDir$(TrustedPathHelper.scala:81) at org.apache.spark.sql.SparkSession.createTrustedTempDir(SparkSession.scala:87) at org.apache.spark.sql.TrustedPathHelper.getTrustedPath(TrustedPathHelper.scala:129) at org.apache.spark.sql.TrustedPathHelper.getTrustedPath$(TrustedPathHelper.scala:128) at org.apache.spark.sql.SparkSession.getTrustedPath(SparkSession.scala:87) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:251) at java.lang.Thread.run(Thread.java:748)
已知API返回内容正常,代码打包为whl在Databricks集群运行,PySpark使用Docker镜像,异常明确发生在json_data = spark.sparkContext.parallelize([response.text])行。
解答
错误原因
该错误与API调用逻辑无关,核心是Spark在创建可信临时目录时的身份验证问题:
- Databricks环境下的
parallelize方法会将本地内存中的数据写入集群节点的可信临时路径,此过程需要验证当前运行用户的系统身份(User Principal) - 由于PySpark运行在Docker镜像中,容器内的用户ID可能在Databricks集群宿主机上不存在,导致系统无法找到对应用户的Principal,从而抛出
UserPrincipalNotFoundException
修复方案
方案1:直接构造DataFrame(推荐)
绕开parallelize的本地文件写入逻辑,将API返回的JSON字符串解析为Python字典后,用spark.createDataFrame直接创建DataFrame:
import json def api_json(spark, api_param): try: token = generate_a_token() token_headers = {'Authorization': f"Bearer {token}"} response = requests.get(f'https://api_url/?api_end_point={api_param}', headers=token_headers) # 解析JSON字符串为Python字典列表 json_dict_list = [json.loads(response.text)] df = spark.createDataFrame(json_dict_list) df.show() return df except Exception as error: traceback.print_exc()
方案2:禁用可信路径检查(仅临时测试)
通过Spark配置关闭可信路径验证,此方法会降低安全性,不建议生产环境使用:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("APIResponseToDF") \ .config("spark.sql.trustedPath.enabled", "false") \ .getOrCreate()
补充:Token生成函数的优化
原generate_a_token函数在异常分支未给token赋值,会导致返回None引发后续API调用失败,建议补充初始化和异常处理:
def generate_a_token(): token = None token_data = { "some_key1": "some_value1", "some_key2": "some_value2" } headers = {"Content-type": "application/x-www-form-urlencoded"} data = bytes(urlencode(token_data).encode()) req = request.Request(url, data, headers) try: response = request.urlopen(req) json_response = json.load(response) token = json_response["access_token"] print(token) except Exception as e: print(f'获取Token失败: {str(e)}') # 可选:抛出异常终止流程,避免传递无效Token # raise e return token
内容的提问来源于stack exchange,提问作者Metadata
相关产品推荐
相关产品推荐

