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

如何使用Spark读取需鉴权HTTPS接口的多JSON响应并转换为DataFrame

实现方案

完全可以实现,根据你的数据量大小可以选择两种不同的实现思路:

方案1:小数据量场景(Driver单请求拉取全量)

如果接口返回的数据量不大,直接在Driver端发起HTTP请求拉取全量数据,再转换为DataFrame即可,步骤如下:

  • 构造HTTP请求,填入你需要的Authorization token、自定义请求头,处理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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 03:54:02