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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 20:05:16