Prefect & Fivetran: Python에서 모든 도구 통합 및 오케스트레이션

Jan 11 2023
간편한 클라우드 데이터 수집 및 오케스트레이션
프리펙트란? 다양한 실행 및 데이터 액세스 패턴을 지원하면서 데이터 흐름을 구축, 안정적으로 실행 및 관찰할 수 있는 유연한 프레임워크입니다. 모든 Python 스크립트를 완전히 운영 가능한 애플리케이션으로 전환할 수 있습니다.
저자의 이미지

프리펙트란?

다양한 실행 및 데이터 액세스 패턴을 지원하면서 데이터 흐름을 구축, 안정적으로 실행 및 관찰할 수 있는 유연한 프레임워크입니다 . 이를 통해 모든 Python 스크립트를 완전히 운영 가능한 애플리케이션 으로 전환할 수 있습니다 .

파이브트란이란?

기업이 다양한 시스템의 데이터를 중앙 데이터 웨어하우스로 복제할 수 있게 해주는 데이터 통합 ​​플랫폼입니다. 데이터베이스, SaaS 애플리케이션, 클라우드 스토리지, 이벤트 로그 등을 포함한 광범위한 데이터 소스를 위한 사전 구축된 커넥터를 제공합니다. Fivetran은 데이터 통합 ​​프로세스를 자동화하여 데이터를 일관되고 완전하며 최신 상태로 유지할 수 있도록 도와줍니다.

왜 Fivetran + Prefect인가?

Prefect는 전체 스택에서 데이터 흐름을 조정합니다. Fivetran을 사용하면 소스에서 대상으로 데이터를 쉽게 이동할 수 있습니다. Prefect를 사용하면 다음을 수행할 수 있습니다.

  • 데이터 복제 동기화를 예약하고 관찰하십시오.
  • 지정된 소스 시스템의 최신 데이터에서 작동해야 하는 파이프라인 부분에 동기화를 추가합니다 .
  • 재시도, 경고 및 종속성 관리를 통해 안정적인 데이터 흐름을 보장합니다.

먼저 무료 Prefect Cloud 계정 에 가입 하고 작업 공간을 만듭니다 . 그런 다음 Prefect를 설치하고 터미널에서 클라우드 작업 영역에 로그인합니다.

pip install prefect
prefect cloud login

그런 다음 다음과 같은 스크립트를 만듭니다 myflow.py.

from prefect import flow

@flow(log_prints=True)
def hello():
    print("You're the Prefectionist now!")

if __name__ == "__main__":
    hello()

Fivetran 시작하기

Fivetran을 시작하려면 평가판 계정에 가입하십시오 . 그런 다음 시작 화면에서 소스, 대상 및 통합할 데이터를 선택하여 첫 번째 Fivetran 동기화를 구성하는 초기 설정을 안내합니다.

소스 연결

시작하려면 GitHub를 커넥터로 선택하고 동기화할 리포지토리를 선택할 수 있습니다. 이렇게 하면 SQL 쿼리를 실행하여 GitHub 별, 풀 요청 등을 분석할 수 있습니다.

목적지 연결

다양한 옵션에서 목적지를 선택할 수 있습니다. 대부분은 클라우드 데이터 웨어하우스 및 클라우드 데이터베이스입니다.

이 데모에서는 Snowflake를 선택하겠습니다. 대상을 선택하면 Snowflake 워크시트에서 실행하여 모든 것을 설정할 수 있는 SQL 트랜잭션 조각이 생성됩니다. 해당 스크립트를 실행하자마자 왼쪽 양식에서 해당 필드( 사용자 이름, 암호, 데이터베이스, 역할 )를 설정할 수 있습니다.

Fivetran 설명서 에서 동일한 코드 블록을 사용할 수 있습니다 . 연결 테스트에 통과하면 데이터 선택을 계속할 수 있습니다.

데이터 선택

개인 정보 보호 및 규정 준수를 위해 일부 열을 차단하거나 해시할 수 있습니다.

이 데모에서는 "모든 데이터 동기화"를 선택하고 "계속"을 클릭합니다.

새 커넥터: Google 스프레드시트

또 다른 간단한 커넥터는 Google 시트에서 데이터를 복제할 수 있습니다 . 다음 표 와 유사한 것을 만들 수 있습니다 .

테이블을 만든 후 데이터 → 명명된 범위 로 이동합니다 .

처음 세 열을 선택하고 이름이 지정된 범위로 저장합니다(예: Fivetran. 예를 들어 에서 열 A, B 및 C의 모든 행을 설정하려면 를 Sheet1사용할 수 있습니다 Sheet1!A:C. 이것은 Fivetran이 동기화하려는 행과 열을 알기 위해 중요합니다.

Fivetran에서 시트 URL을 붙여넣은 다음 명명된 범위를 선택합니다. 그런 다음 커넥터를 저장하고 테스트합니다.

모든 것이 예상대로 작동하면 "모든 데이터 동기화" 및 "계속"을 선택할 수 있습니다.

이제 초기 동기화를 시작할 수 있습니다.

동기화 후에는 Snowflake 워크시트에 해당 Google 시트에 해당하는 새 스키마와 테이블이 표시되어야 합니다.

Prefect & Fivetran으로 데이터 복제 자동화

Fivetran API 키 생성

동기화를 예약하고 조정하려면 Prefect 흐름이 Fivetran에 인증되어야 합니다. 이를 위해서는 API 키가 필요합니다. 계정 설정으로 이동하여 아래와 같이 API 키를 생성합니다.

이렇게 하면 API 키 식별자와 API 비밀 토큰이 제공됩니다 .

Fivetran을 위한 Prefect 블록 만들기

prefect-fivetran 컬렉션을 설치 하고 Fivetran 블록을 등록합니다.

pip install prefect prefect_fivetran
prefect block register -m prefect_fivetran

코드로 이를 수행하는 방법은 다음과 같습니다.

from dotenv import load_dotenv
import os
from prefect_fivetran import FivetranCredentials

load_dotenv()

fivetran_credentials = FivetranCredentials(
    api_key=os.environ.get("FIVETRAN_API_KEY"),
    api_secret=os.environ.get("FIVETRAN_API_SECRET_KEY"),
)
fivetran_credentials.save("default")

마지막으로 흐름을 실행하여 gsheet.py동기화를 트리거할 수 있습니다.

from prefect import flow
from prefect_fivetran import FivetranCredentials
from prefect_fivetran.connectors import (
    wait_for_fivetran_connector_sync,
    start_fivetran_connector_sync,
)


@flow
def example_flow(connector_id: str):
    fivetran_credentials = FivetranCredentials.load("default")

    last_sync = start_fivetran_connector_sync(
        connector_id=connector_id,
        fivetran_credentials=fivetran_credentials,
    )

    return wait_for_fivetran_connector_sync(
        connector_id=connector_id,
        fivetran_credentials=fivetran_credentials,
        previous_completed_at=last_sync,
        poll_status_every_n_seconds=60,
    )


if __name__ == "__main__":
    example_flow("bereft_indices")

데이터 복제 동기화 트리거

Fivetran, Prefect 및 Snowflake 통합이 제대로 설정되었는지 확인하기 위해 이 흐름을 두 번 트리거합니다. 먼저 Google 시트 테이블을 있는 그대로 사용합니다. 그런 다음 일부 레코드를 수정하고 동기화 워크플로가 대상에서 레코드를 업데이트했는지 검사합니다. 이렇게 하면 Snowflake 테이블이 소스 시스템과 동기화 상태를 유지하는지 확인할 수 있습니다.

이 흐름을 처음으로 트리거해 보겠습니다.

python gsheet.py

      
                

동기화에 성공했으며 첫 번째 흐름 실행의 결과로 Snowflake 테이블에서 변경된 사항이 없습니다.

이제 첫 번째 고객의 이름을 Michael에서 Mike로 변경하고 ID가 2인 Shawn의 고객 레코드를 제거하고 ID가 101인 다른 고객을 추가해 보겠습니다.

변경한 후 Prefect 흐름을 다시 트리거합니다.

python gsheet.py

수정된 데이터가 Snowflake 테이블에 올바르게 복제되었는지 확인할 수 있습니다.

천천히 변화하는 치수

Michael에서 Mike로의 이름 업데이트와 Shawn의 레코드 삭제를 반영하는 변경 사항은 이 테이블에 표시되지 않습니다. 느린 변경 차원 전략 의 일부로 해당 정보를 추적하려면 해당 동기화에 대한 기록을 활성화해야 합니다 . 추적 기록을 지원하는 커넥터 목록은 기록 모드 Fivetran 문서를 확인하세요 .

동기화 예약

이 스크립트를 일정에 따라(예: 매일 오전 9시) 실행하려는 경우 다음 명령을 사용하여 Prefect 배포를 생성할 수 있습니다 .

prefect deployment build --cron "0 9 * * *" gsheet.py:example_flow -n dev -a

prefect agent start -q default

다음 단계

이 게시물은 Prefect 및 Fivetran으로 시작하여 GitHub 및 Google Sheets에서 Snowflake로 첫 번째 데이터 복제 동기화를 생성하는 방법을 보여주었습니다. Prefect 블록을 사용하여 API 키를 안전하게 저장하고 Prefect 흐름에서 Fivetran 동기화를 오케스트레이션하는 방법을 살펴보았습니다.

이 게시물에서 논의된 내용이 명확하지 않은 경우 Prefect Community Slack 또는 Prefect Discourse 에서 질문할 때 저를 태그해 주세요 .

읽어주셔서 감사합니다. 즐거운 엔지니어링 되세요!