PySpark UDF如何同时返回状态码与响应值并存入DataFrame不同列
实现方案
你可以通过让PySpark UDF返回结构化类型StructType的方式实现多值返回,再拆分到不同列即可,具体实现如下:
步骤1:修正原有UDF逻辑 + 调整返回值
首先修正原有代码的逻辑错误(原代码提前将response对象转为JSON覆盖了原变量,无法再取status_code),同时让函数返回状态码和响应结果的组合:
import requests import json from pyspark.sql.functions import udf, col from pyspark.sql.types import StructType, StructField, IntegerType, MapType, StringType def Api(a): path = endpoint headers = {'sample-Key': sample} # 这里原代码body写死,替换为传入的参数a body = [{'text': a}] status_code = None output = None try: # 先保留原始response对象,不要直接转JSON覆盖变量 resp = requests.post(path, params=params, headers=headers, json=body) status_code = resp.status_code if status_code == 200: # 如果需要输出JSON字符串,替换为 json.dumps(resp.json()) output = resp.json() except Exception as e: # 异常场景自定义状态码和输出内容 status_code = 500 output = str(e) # 同时返回两个值 return (status_code, output)
步骤2:定义UDF的结构化返回类型
根据你返回的output格式调整Schema,示例如下:
# 若output是JSON对象则用MapType,若为字符串则替换为StringType() udf_schema = StructType([ StructField("status_code", IntegerType(), nullable=True), StructField("output", MapType(StringType(), StringType()), nullable=True) ]) udf_Api = udf(Api, udf_schema)
步骤3:调用UDF并拆分到不同列
优先选择以下方式,仅调用一次UDF,性能更高:
# 先生成中间结构化列,再拆分出目标字段 temp_df = df.withColumn("api_res", udf_Api(col("input"))) new_df = temp_df.select( "input", col("api_res.status_code").alias("status_code"), col("api_res.output").alias("output") ) # 不需要保留中间列的话可以执行 temp_df.drop("api_res")
注意:不要直接在select中两次调用UDF取不同字段,会导致每个行触发两次API请求,性能损耗极大
内容的提问来源于stack exchange,提问作者dcrowley01
相关产品推荐
相关产品推荐

