如何使用Spark读取需鉴权HTTPS接口的多JSON响应并转换为DataFrame
实现方案
完全可以实现,根据你的数据量大小可以选择两种不同的实现思路:
方案1:小数据量场景(Driver单请求拉取全量)
如果接口返回的数据量不大,直接在Driver端发起HTTP请求拉取全量数据,再转换为DataFrame即可,步骤如下:
- 构造HTTP请求,填入你需要的
Authorizationtoken、自定义请求头,处理HTTPS验证(公网通用证书默认信任即可,自定义证书可按需配置) - 解析返回的响应体,按分隔规则拆分出所有独立的JSON字符串
- 将JSON字符串列表转为RDD,直接调用Spark的JSON解析接口生成DataFrame
Python PySpark示例代码:
import requests from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ReadHttpJson").getOrCreate() # 配置请求参数 url = "你的HTTPS接口地址" headers = { "Authorization": "Bearer 你的token", "Content-Type": "application/json", # 补充其他需要的请求头 } # 发起请求 response = requests.get(url, headers=headers, verify=True) response.raise_for_status() # 拆分多个独立JSON,此处按换行拆分,若返回的JSON无换行可按正则匹配提取所有独立JSON对象 json_str_list = [line.strip() for line in response.text.splitlines() if line.strip().startswith("{")] # 转RDD后解析为DataFrame json_rdd = spark.sparkContext.parallelize(json_str_list) df = spark.read.json(json_rdd) # 验证结果 df.show()
方案2:大数据量场景(分布式并行拉取)
如果接口支持分页、分片拉取,数据量较大的情况下可以用分布式并行拉取避免Driver内存溢出:
- 先在Driver端生成分页参数列表(比如页码、分片ID列表)
- 将分页参数转为RDD,使用
mapPartitions算子在每个Executor节点上发起HTTP请求拉取对应分片的数据 - 每个Executor内解析JSON字符串后返回,最后统一合并为DataFrame
Scala示例代码:
import org.apache.spark.sql.SparkSession import sttp.client3._ object ReadHttpJson { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("ReadHttpJson").getOrCreate() import spark.implicits._ // 生成分页参数列表,示例为1到100页 val pageList = (1 to 100).toList val pageRdd = spark.sparkContext.parallelize(pageList, 10) // 配置10个并行拉取的分区 val jsonRdd = pageRdd.mapPartitions(iter => { // 每个分区复用HTTP客户端,减少连接开销 val backend = HttpURLConnectionBackend() iter.flatMap(page => { val response = basicRequest .get(uri"你的HTTPS地址?page=$page") .header("Authorization", "Bearer 你的token") // 补充其他需要的请求头 .send(backend) response.body match { case Right(body) => body.splitlines().filter(_.trim.startsWith("{")) case Left(_) => List.empty[String] } }) }) val df = spark.read.json(jsonRdd.toDS()) df.show() } }
关键注意事项
- 若返回的多个JSON没有换行分隔,可以用正则匹配
\{.*?\}提取所有独立JSON对象,注意处理JSON内部嵌套大括号的特殊情况 - 高并发拉取时建议添加限流逻辑,避免请求频率过高被接口方拦截
- 如果是自签名HTTPS证书,需要在HTTP客户端中配置信任证书,生产环境不推荐临时关闭证书验证
内容的提问来源于stack exchange,提问作者Abhishek kumar
相关产品推荐
相关产品推荐

