PySpark调用REST API打印结果报collapse_columns未定义的解决方法
报错原因
运行抛出name 'collapse_columns' is not defined是因为collapse_columns不是PySpark官方提供的内置函数,是你参考的技术文章作者自行编写的嵌套结构展平工具函数,复制代码时漏掉了该函数的定义,导致运行时找不到对应方法。
直接打印DataFrame只能看到简略嵌套输出,是因为explode处理后得到的results列是Struct结构体类型,包含Make_ID、Make_Name两个子字段,未展开结构体的情况下,show()方法只会打印结构体的简略占位符,不会展示内部具体字段值。
修复步骤
- 补全缺失的
collapse_columns自定义函数,该函数作用是递归遍历Schema,把嵌套的Struct字段展平为顶层单列,无需手动逐个声明子字段选择逻辑 - (可选)修正原有REST请求函数的逻辑问题:原函数入参接收了自定义headers,但函数内部硬编码覆盖了headers变量,会导致外部传入的请求头配置不生效
- (可选简化方案)如果返回结构固定且简单,也可以不使用自定义展平函数,直接选择结构体下的子字段即可
补全的工具函数代码
from pyspark.sql.functions import col from pyspark.sql.types import StructType def collapse_columns(schema, prefix=None): field_list = [] for field in schema.fields: full_col_name = f"{prefix}.{field.name}" if prefix else field.name # 遇到Struct类型则递归展开 if isinstance(field.dataType, StructType): field_list.extend(collapse_columns(field.dataType, full_col_name)) else: # 展开后的字段名用下划线替代点号,避免嵌套字段名歧义 field_list.append(col(full_col_name).alias(full_col_name.replace(".", "_"))) return field_list
简化写法(无需自定义函数)
当前接口返回的Results数组展开后只有两个子字段,可以直接手动选择字段,不需要额外写展平函数:
df = result_df.select(explode(col("result.Results")).alias("results")) # 直接选择结构体下的子字段 df.select( col("results.Make_ID").alias("Make_ID"), col("results.Make_Name").alias("Make_Name") ).show(truncate=False)
修正后的REST请求函数(可选优化)
def executeRestApi(verb, url, headers, body): res = None try: if verb == "get": res = requests.get(url, data=body, headers=headers) else: res = requests.post(url, data=body, headers=headers) except Exception as e: return e if res is not None and res.status_code == 200: return json.loads(res.text) return None
运行结果
修复后执行show()即可正常展示解析后的完整JSON数据,不再显示[,,,]这类简略占位输出,结果格式如下:
+-------+-------------------+ |Make_ID|Make_Name | +-------+-------------------+ |440 |ASTON MARTIN | |441 |TESLA | |442 |JAGUAR | |443 |LAND ROVER | |... |... | +-------+-------------------+
内容的提问来源于stack exchange,提问作者Ana
相关产品推荐
相关产品推荐

