在管道中从API中获取数据

从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。 由于该数据集是物化视图,因此每次管道更新时,都会以完整且幂等的方式重新运行该函数。

以下步骤教你如何构建带有周期性拉取的具体化视图:

  1. 把 API 令牌存入一个秘密,然后映射到流水线设置中的 Spark 配置属性,这样流水线代码就能读取它。 将该属性添加到 spark_conf 管道集群配置的块中:

    {
      "clusters": [
        {
          "spark_conf": {
            "api.token": "{{secrets/<scope-name>/<secret-name>}}"
          }
        }
      ]
    }
    

    下一步中的代码使用 spark.conf.get("api.token") 读取该值。 关于如何在管道设置中配置秘密的更多信息,请参见 “安全访问带有管道内秘密的存储凭证”。

  2. 定义一个物质化视图,调用 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)
    
  3. 在函数内部处理分页,方法是循环页面并串联结果,然后再返回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 拉取数据。

以下步骤会向您展示如何从自定义数据源中导入:

  1. 实现一个 DataSourceDataSourceStreamReader,用于调用 API 并跟踪读取偏移量。 有关创建自定义数据源的详细信息,请参见 PySpark 自定义数据源

  2. 注册该数据源,以便管道能够按格式名称引用它:

    spark.dataSource.register(MyApiDataSource)
    
  3. 从流式处理表中的已注册的源读取:

    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 特有的兼容性问题(如分页和速率限制)与声明式转换逻辑隔离开来,同时免费获得自动加载器的精确一次文件跟踪能力。

以下步骤展示了如何将摄入与计划工作分离:

  1. 写一个笔记本或脚本,调用 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)
    
  2. 通过Lakeflow Jobs安排笔记本或脚本独立运行。 请参阅 Lakeflow Jobs

  3. 在你的流水线中,定义一个流式表,用自动加载器读取已着陆的文件:

    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 响应流入下游之前将其捕获。
  • 处理分页和速率限制。 遍历各个页面,并添加带退避的重试机制,以免瞬时故障导致整个更新失败。

其他资源