如何通过PySpark/Databricks直接将API返回的JSON读取为DataFrame?
解决方案:跳过文件直接将API响应转为Databricks DataFrame
你的代码存在两个核心问题:
- 误用
json.dumps(response.text):response.text本身就是API返回的合法JSON字符串,json.dumps会对其进行二次转义(例如把{"id":1}变成"{\"id\":1}"),导致Spark无法解析为正常JSON结构,最终生成_corrupt_record列。 - 直接传递字符串给
spark.read.json():该方法默认接受文件路径参数,而非原始JSON字符串,因此会触发"relative path in URI"错误。
以下是针对不同场景的可行解决方案:
方案1:直接将API响应转为Python对象创建DataFrame
这是最简洁的方式,适用于API返回单个JSON对象或JSON数组的场景:
import requests from pyspark.sql import SparkSession # 初始化SparkSession(若全局已初始化可省略) spark = SparkSession.builder.appName("APItoDF").getOrCreate() payload = {} headers = { 'Authorization': 'Basic ==', 'Cookie': 'JSESSIONID=' } response = requests.request("GET", apipath, headers=headers, data=payload) # 将响应直接转为Python字典/列表 data = response.json() # 处理单个对象或数组的情况 if isinstance(data, dict): df = spark.createDataFrame([data]) # 单个对象需包装为列表 else: df = spark.createDataFrame(data) # 数组可直接传入 df.printSchema() df.show()
方案2:使用RDD传递原始JSON字符串给spark.read.json
如果需要保留spark.read.json的使用方式,需确保传递的是未转义的原始JSON字符串RDD:
import requests from pyspark.sql import SparkSession spark = SparkSession.builder.appName("APItoDF").getOrCreate() payload = {} headers = { 'Authorization': 'Basic ==', 'Cookie': 'JSESSIONID=' } response = requests.request("GET", apipath, headers=headers, data=payload) # 直接获取原始JSON字符串,避免二次转义 raw_json = response.text # 将字符串转为RDD json_rdd = spark.sparkContext.parallelize([raw_json]) # 读取RDD生成DataFrame df = spark.read.json(json_rdd) df.printSchema() df.show()
方案3:循环调用API的批量处理
针对循环调用多API的场景,可收集所有响应数据后统一生成DataFrame:
import requests from pyspark.sql import SparkSession spark = SparkSession.builder.appName("BatchAPItoDF").getOrCreate() payload = {} headers = { 'Authorization': 'Basic ==', 'Cookie': 'JSESSIONID=' } # 示例:多个API路径列表 api_paths = ["api/path/1", "api/path/2", "api/path/3"] all_response_data = [] for apipath in api_paths: response = requests.request("GET", apipath, headers=headers, data=payload) data = response.json() # 区分单个对象和数组,统一收集到列表 if isinstance(data, dict): all_response_data.append(data) else: all_response_data.extend(data) # 批量生成DataFrame df = spark.createDataFrame(all_response_data) df.printSchema() df.show()
内容的提问来源于stack exchange,提问作者Lynchie
相关产品推荐
相关产品推荐

