파이프라인에서 API의 데이터를 수집하기

API에서 인제스트한다는 것은 파일이나 데이터베이스에서 읽는 대신 웹 서비스에서 보통 페이지네이트 JSON 형태로 데이터를 HTTP로 가져오는 것을 의미합니다. 파일이나 메시지 버스와 달리 내장 일반 API 소스가 없기 때문에 인증, 페이지닝, 속도 제한은 직접 처리해야 합니다. Lakeflow 파이프라인은 임의의 API에서 인제스트하는 세 가지 패턴을 지원합니다. 어느 것이 적합한지는 처리량과 새로 고침 요구 사항에 따라 달라집니다.

Important

커스텀 API 인제스 코드를 작성하기 전에, 소스에 이미 관리되는 커넥터가 있는지 확인하세요. Lakeflow Connect는 Salesforce, Workday, ServiceNow, Google Analytics 등 다양한 일반적인 소프트웨어-서비스(SaaS) API용 내장 커넥터를 제공하며, 파트너 커넥터도 점점 더 많이 제공되고 있습니다. 커넥터가 소스를 덮으면 인증, 페이지네이션, 점진적 추출을 대신 처리하며, 거의 항상 손으로 말아 입력하는 것보다 훨씬 덜 번거롭습니다. Lakeflow Connect 커넥터 개념을 참조하세요. 아래 패턴은 커넥터가 맞지 않을 때만 사용하세요.

사전 요구 사항

  • 파이프라인. 하나 만들려면 Lakeflow 파이프라인 튜토리얼을 참고하세요.
  • 토큰이나 키와 같은 API 자격 증명은 Azure Databricks 비밀로 저장됩니다. 파이프라인 소스 코드에 자격 증명을 하드코딩하지 마세요. 비밀 관리를 참조하세요.
  • 파이프라인 컴퓨트에서 API 엔드포인트로 네트워크 접근을 합니다.
  • 스트리밍 테이블과 물질화된 뷰에 대한 익숙함, 이러한 패턴이 생성하는 데이터셋 유형에 대한 이해. 스트리밍 테이블물질화된 뷰를 참조하세요.

패턴 선택

파이프라인에는 네이티브한 일반 REST-API 소스가 없으므로, 임의의 API에서 데이터를 가져올 때는 데이터 양과 인징 빈도에 따라 세 가지 패턴 중 하나를 선택하세요:

패턴 사용 시기
구체화된 뷰로서의 주기적 가져오기 페이로드는 소형~중형 규모이며, 참조 데이터, 일일 외환 환율, 또는 페이지네이션이 적용되지만 범위를 제한할 수 있는 API처럼 파이프라인 실행당 한 번만 가져옵니다.
Python 데이터 소스 API 대용량 또는 스트리밍 API를 점진적으로 폴링하고, 체크포인트로 진행 상태를 기록해 재시작 시 전체를 다시 읽지 않도록 해야 합니다.
Auto Loader를 사용한 분리된 데이터 수집 API별 특이점을 변환 로직에서 분리하고, 별도 작업 없이 파일을 정확히 한 번만 추적할 수 있기를 원합니다.

패턴 1: 물질화된 시각으로서의 주기적 당김

파이프라인 실행당 한 번씩 소규모에서 중간 규모의 페이로드를 작성할 때, API를 호출하고 Spark DataFrame을 반환하는 Python 함수를 작성하세요. 데이터셋이 물질화된 뷰이기 때문에, 파이프라인은 업데이트할 때마다 함수를 완전하고 면치있게 다시 실행합니다.

다음 단계들은 주기적인 당김이 있는 구체화된 뷰를 만드는 방법을 보여줍니다:

  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를 사용하세요.

패턴 2: Python 데이터 소스 API를 사용하는 대용량 또는 스트리밍 API

API는 오프셋 트래킹으로 점진적으로 폴링해야 하며, Spark의 Python Data Source API를 사용해 맞춤형 데이터 소스를 구현하세요. 이렇게 하면 체크포인트 진행 상황과 점진적 읽기 등 적절한 스트리밍 의미론이 제공되어, 재시작이 API 전체를 다시 가져오는 대신 마지막 오프셋부터 다시 시작됩니다.

다음 단계들은 사용자 지정 데이터 소스에서 인제스트하는 방법을 보여줍니다:

  1. API를 호출하는 a DataSourceDataSourceStreamReader 를 구현하고 읽기 오프셋을 추적하세요. 커스텀 데이터 소스 작성에 대한 자세한 내용은 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()
    

패턴 3: 예약 작업과 Auto Loader를 사용해 수집 분리

일반적인 생산 패턴은 API 호출을 파이프라인과 분리하는 것입니다. 예약된 작업은 원시 API 응답을 Unity 카탈로그 볼륨에 파일로 저장하고, 파이프라인이 Auto Loader로 이를 처리합니다. 이 기능은 페이지네이션이나 속도 제한 같은 API별 특이점들을 선언적 변환 로직에서 분리하고, Auto Loader의 정확히 한 번만 가능한 파일 추적을 무료로 제공합니다.

다음 단계들은 정해진 작업과 섭취를 분리하는 방법을 보여줍니다:

  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 작업을 참조하세요.

  3. 파이프라인에서 Auto Loader로 접지된 파일을 읽는 스트리밍 테이블을 정의하세요:

    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")
        )
    

Auto Loader를 이용한 신뢰할 수 있는 파일 인제스팅에 대해 더 알고 싶다면, 'Cloud Object 저장소에서 파일 로드 '와 ' Auto Loader란 무엇인가?'를 참고하세요.

API 인제스팅의 모범 사례

  • 소스 코드에 비밀을 넣지 마세요. API 토큰과 키를 Azure Databricks의 비밀 범위에 저장하고 런타임에 읽으세요. 비밀 관리를 참조하세요.
  • 답변을 일찍부터 검증하세요. 인제스트된 행에 기대 치를 추가하여 잘못된 API 응답이 하류로 흐르기 전에 포착하세요.
  • 페이지네이션과 속도 제한을 관리하세요. 페이지를 루프 오버하고 백오프와 함께 재시도를 추가해서 일시적 실패가 전체 업데이트를 실패하지 않도록 하세요.

추가 리소스