从API中获取数据意味着通过HTTP从Web服务中拉取数据,通常是分页的JSON,而不是从文件或数据库读取。 与文件或消息总线不同,它没有内置通用的 API 源,所以认证、分页和速率限制都是自己处理的。 Lakeflow 管道支持三种从任意 API 引入数据的模式。 哪个适合取决于你的音量和刷新率需求。
重要
在编写任何自定义 API 导入代码之前,先检查一下你的源代码是否已经有托管连接器。 Lakeflow Connect 内置连接器支持许多常见的软件即服务(SaaS)API,如 Salesforce、Workday、ServiceNow 和 Google Analytics,合作伙伴连接器数量也在不断增加。 如果连接器覆盖了你的源,它会帮你处理身份验证、分页和增量提取,而且几乎总比手动编写数据引入更省事。 参见 Lakeflow Connect连接器概念。 仅当没有适用的连接器时,才使用以下模式。
先决条件
- 一条管道。 要创建一个,请参见 Lakeflow管道教程。
- API 凭证,如令牌或密钥,作为 Azure Databricks 秘密存储。 绝不要在流水线源代码中硬编码凭证。 请参阅机密管理。
- 从你的管道计算资源到 API 终结点的网络访问。
- 熟悉流式表和物化视图,以及这些模式生成的数据集类型。 参见 流式表 和 实体化视图。
选择一个图案
流水线中没有原生的通用 REST-API 源,所以当你从任意API中提取数据时,根据数据量和摄入频率选择三种模式之一:
| 图案 | 何时使用 |
|---|---|
| 周期性拉取作为具体化视图 | 有效负载的规模从较小到中等,并且在每次管道运行时拉取一次,例如参考数据、每日外汇汇率,或分页但范围可限定的 API。 |
| Python 数据源 API | 你需要以增量方式轮询高吞吐量或流式 API,并记录检查点进度,这样在重启后就不必从头重新读取全部数据。 |
| 使用 Auto Loader 实现解耦式数据摄取 | 你希望将 API 特有的兼容性问题与数据转换逻辑分离开来,免费获得精确一次的文件跟踪。 |
模式 1:周期性拉取作为具体化视图
对于每次流水线拉取一次的中小型负载,编写一个 Python 函数调用 API,返回 Spark DataFrame。 由于该数据集是物化视图,因此每次管道更新时,都会以完整且幂等的方式重新运行该函数。
以下步骤教你如何构建带有周期性拉取的具体化视图:
把 API 令牌存入一个秘密,然后映射到流水线设置中的 Spark 配置属性,这样流水线代码就能读取它。 将该属性添加到
spark_conf管道集群配置的块中:{ "clusters": [ { "spark_conf": { "api.token": "{{secrets/<scope-name>/<secret-name>}}" } } ] }下一步中的代码使用
spark.conf.get("api.token")读取该值。 关于如何在管道设置中配置秘密的更多信息,请参见 “安全访问带有管道内秘密的存储凭证”。定义一个物质化视图,调用 API,并以 DataFrame 返回响应:
import requests from pyspark import pipelines as dp from pyspark.sql import Row @dp.materialized_view( name="exchange_rates_bronze", comment="Daily FX rates pulled from a public REST API", ) def exchange_rates_bronze(): resp = requests.get( "https://api.example.com/v1/rates", params={"base": "USD"}, headers={"Authorization": f"Bearer {spark.conf.get('api.token')}"}, timeout=30, ) resp.raise_for_status() rates = resp.json()["rates"] rows = [Row(currency=k, rate=float(v), as_of_date=resp.json()["date"]) for k, v in rates.items()] return spark.createDataFrame(rows)在函数内部处理分页,方法是循环页面并串联结果,然后再返回DataFrame:
import requests from pyspark import pipelines as dp from pyspark.sql import Row @dp.materialized_view( name="customers_bronze", comment="Customers pulled from a paginated REST API", ) def customers_bronze(): token = spark.conf.get("api.token") rows = [] url = "https://api.example.com/v1/customers" while url: # follow the API's next-page cursor until exhausted resp = requests.get( url, headers={"Authorization": f"Bearer {token}"}, timeout=30, ) resp.raise_for_status() payload = resp.json() rows.extend(Row(**record) for record in payload["data"]) url = payload.get("next") # None on the last page return spark.createDataFrame(rows)为该请求添加重试和退避逻辑,以提高系统韧性。
该模式在每次管道更新时重新读取完整的 API 响应,因此仅在有效载荷有边界时使用。 对于增量读取,使用模式2。
模式二:使用Python数据源API进行高流量或流式API
对于需要通过偏移量跟踪进行增量轮询的 API,可以使用 Spark 的 Python Data Source API 来实现自定义数据源。 这为你提供了正确的流式处理语义,包括带检查点的进度记录和增量读取,因此重启时会从上次的偏移量继续恢复,而不是再次从整个 API 拉取数据。
以下步骤会向您展示如何从自定义数据源中导入:
实现一个
DataSource和DataSourceStreamReader,用于调用 API 并跟踪读取偏移量。 有关创建自定义数据源的详细信息,请参见 PySpark 自定义数据源。注册该数据源,以便管道能够按格式名称引用它:
spark.dataSource.register(MyApiDataSource)从流式处理表中的已注册的源读取:
from pyspark import pipelines as dp @dp.table(name="events_bronze") def events_bronze(): return spark.readStream.format("my_api_source").load()
模式三:将导入与定时作业和自动加载器分离
一个常见的生产模式是将API调用与流水线分离。 调度作业会将原始 API 响应以文件形式写入 Unity Catalog 卷中,而该流水线则通过 Auto Loader 读取这些文件。 这样可以将 API 特有的兼容性问题(如分页和速率限制)与声明式转换逻辑隔离开来,同时免费获得自动加载器的精确一次文件跟踪能力。
以下步骤展示了如何将摄入与计划工作分离:
写一个笔记本或脚本,调用 API,并向 Unity 目录卷写出原始的 JSON 响应。 从机密中读取 API 凭证。 请参阅机密管理。
import requests, json, time token = dbutils.secrets.get(scope="<scope-name>", key="<secret-name>") volume_path = "/Volumes/main/raw/landing/api_events" resp = requests.get( "https://api.example.com/v1/events", headers={"Authorization": f"Bearer {token}"}, timeout=30, ) resp.raise_for_status() # One file per run; the pipeline's Auto Loader tracks which files it has ingested. with open(f"{volume_path}/events_{int(time.time())}.json", "w") as f: json.dump(resp.json()["data"], f)通过Lakeflow Jobs安排笔记本或脚本独立运行。 请参阅 Lakeflow Jobs。
在你的流水线中,定义一个流式表,用自动加载器读取已着陆的文件:
from pyspark import pipelines as dp @dp.table(name="api_events_bronze") def api_events_bronze(): return ( spark.readStream.format("cloudFiles") .option("cloudFiles.format", "json") .load("/Volumes/main/raw/landing/api_events") )
关于使用自动加载器实现可靠文件摄取的更多信息,请参见 “从云对象存储加载文件 ”和 “什么是自动加载器?”。
API 数据摄取最佳实践
- 把秘密藏在源代码之外。 将 API 令牌和密钥存储在 Azure Databricks 的秘密作用域中,并在运行时读取它们。 请参阅机密管理。
- 尽早验证回复。 为摄取的数据行添加校验规则,以便在格式异常的 API 响应流入下游之前将其捕获。
- 处理分页和速率限制。 遍历各个页面,并添加带退避的重试机制,以免瞬时故障导致整个更新失败。