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
相关产品推荐
相关产品推荐

