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

Spark新手求助:分页REST API数据写入Databricks分区表

分步解决分页REST API数据采集与Databricks分区存储问题

嘿,作为Spark新手碰到分页API的问题真的很常见,我来一步步给你拆解怎么实现从分页REST API取数,再按日期分区存到Databricks表~

1. 实现分页API的循环调用

首先得搞定分页逻辑——既然要调用5次左右,我们可以写个循环来遍历所有页码,直到获取完所有数据。这里以Python为例(Databricks里Python和Scala都常用,选你顺手的就行):

import requests
from pyspark.sql import SparkSession
from pyspark.sql.functions import current_date

# Databricks环境里其实不需要手动初始化SparkSession,不过写出来方便理解
spark = SparkSession.builder.appName("PagedRESTAPIIngest").getOrCreate()

# API基础配置,替换成你的实际信息
api_base_url = "https://your-api-domain.com/your-endpoint"
request_headers = {
    "Authorization": f"Bearer {dbutils.secrets.get(scope='your-secret-scope', key='api-access-token')}",
    "Content-Type": "application/json"
}
page_size = 200  # 按API支持的每页条数调整
current_page = 1
all_records = []

while True:
    # 构造分页请求参数
    request_params = {"page": current_page, "pageSize": page_size}
    # 发送请求
    response = requests.get(api_base_url, headers=request_headers, params=request_params)
    # 主动抛出请求错误,方便排查问题
    response.raise_for_status()
    # 解析返回的JSON
    api_response = response.json()
    
    # 提取当前页的数据(这里假设API把数据放在"data"字段里,根据实际结构调整)
    current_page_records = api_response.get("data", [])
    if not current_page_records:
        # 当前页没数据,说明已经取完了
        break
    
    all_records.extend(current_page_records)
    
    # 判断是否还有下一页,常见的判断方式有两种:
    # 方式1:API返回hasNext标志
    has_next_page = api_response.get("hasNext", False)
    # 方式2:对比当前页码和总页码
    total_pages = api_response.get("totalPages", 1)
    
    if not has_next_page or current_page >= total_pages:
        break
    
    current_page += 1

2. 转换为Spark DataFrame并添加分区日期列

拿到所有数据后,我们把它转换成Spark DataFrame,再添加一个ingest_date列作为分区键——因为你是每日调用,用当日日期就刚好:

# 将收集到的列表转为Spark DataFrame
raw_df = spark.createDataFrame(all_records)

# 添加当日日期作为分区列(current_date()会返回运行脚本当天的日期)
df_with_partition = raw_df.withColumn("ingest_date", current_date())

3. 按日期分区写入Databricks表

最后一步就是把数据写入Databricks表,这里推荐用Delta Lake格式(Databricks官方主推,支持ACID、增量更新等特性),并且指定按ingest_date分区:

# 替换成你的Catalog、Schema和表名
target_table = "your_catalog.your_schema.your_target_table"

# 写入表:用append模式(因为每日跑,要累加数据),按ingest_date分区
df_with_partition.write \
    .mode("append") \
    .partitionBy("ingest_date") \
    .format("delta") \
    .saveAsTable(target_table)

新手友好的额外提示

  • 安全存储密钥:千万别把API密钥硬编码在代码里,用Databricks的Secret Scope来管理,就像示例里的dbutils.secrets.get()那样
  • 异常重试:如果API偶尔会超时,可以加个重试机制,比如用tenacity库来实现自动重试
  • 并行优化:如果API允许并发请求,你可以把要请求的页码列表并行化,用Spark的parallelize来同时请求多个页,加快取数速度
  • 数据校验:写入前可以加一些简单的校验,比如检查关键字段是否为空,避免脏数据入库

这样你就可以把这段代码做成Databricks Job,设置每日调度,自动完成分页API取数和分区存储啦~如果你的API分页逻辑和示例不一样,只要调整循环里的判断条件(比如有的API用offset代替page)就行。

内容的提问来源于stack exchange,提问作者Swati Patil

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 06:30:44