Lakeflow 파이프라인이란?

Lakeflow 파이프라인은 SQL 및 Python 일괄 처리 및 스트리밍 데이터 파이프라인을 빌드하기 위한 선언적 프레임워크를 제공합니다. 이들의 핵심 개념은 파이프라인, 플로우, 스트리밍 테이블, 머터리얼라이즈드 뷰, 싱크이며, 이들은 서로 함께 작동해 자동 오케스트레이션과 증분 업데이트를 통해 데이터를 처리합니다.

Lakeflow 파이프라인은 Apache Spark™ SDP(선언적 파이프라인)를 확장합니다. SDP 및 Lakeflow 파이프라인과 비교하는 방법에 대한 자세한 내용은 Apache Spark 선언적 파이프라인을 참조하세요.

Tip

파이프라인이 처음이신가요? Lakeflow 파이프라인 사용법부터 시작하여 파이프라인을 어떻게 그리고 왜 사용하는지 이해하고, 각 단계별 작업 링크를 확인하세요.

Note

Lakeflow 파이프라인에는 프리미엄 계획이 필요합니다. 자세한 내용은 Databricks 계정 팀에 문의하세요.

파이프라인의 이점은 무엇인가요?

Apache SparkSpark Structured Streaming API를 사용해 Databricks Runtime에서 데이터 엔지니어링 프로세스를 개발할 때 Lakeflow Jobs를 통해 수동으로 오케스트레이션하는 방식과 달리, 파이프라인의 선언적 특성은 다음과 같은 이점을 제공합니다.

  • 자동 오케스트레이션: 파이프라인은 최대 병렬 처리를 사용하여 올바른 순서로 처리 단계("흐름"이라고 함)를 실행하고, Spark 태스크에서 흐름, 전체 파이프라인까지 임시 오류를 점진적으로 다시 시도합니다.
  • 선언적 처리: 선언적 함수는 수백 줄의 수동 Spark 및 구조적 스트리밍 코드를 몇 줄로 줄입니다. AUTO CDC API는 순서가 다른 이벤트 또는 워터마크와 같은 스트리밍 개념에 대한 수동 코드 없이 SCD Type 1 및 Type 2를 포함한 CDC(변경 데이터 캡처) 이벤트를 처리합니다.
  • 증분 처리: 증분 처리 엔진은 구체화된 뷰를 최신 상태로 유지합니다. 일괄 처리 의미 체계를 사용하여 변환 논리를 작성하고 가능한 경우 엔진은 새 데이터 또는 변경된 원본 데이터만 다시 처리합니다.

주요 개념

아래 다이어그램은 파이프라인의 가장 중요한 개념을 보여 줍니다.

파이프라인의 핵심 개념이 매우 높은 수준에서 서로 어떻게 관련되는지 보여 주는 다이어그램

데이터 세트

파이프라인은 각각 다른 처리 의미 체계를 가진 세 가지 유형의 데이터 세트를 생성합니다.

데이터 세트 형식 레코드 처리 방법
스트리밍 테이블 각 레코드는 추가 전용 원본을 가정하여 정확히 한 번 처리됩니다. 스트리밍 테이블은 지속적으로 증가하는 데이터의 수집 및 증분 처리에 적합합니다.
구체화된 뷰 결과는 데이터의 현재 상태를 반영하기 위해 필요에 따라 다시 계산됩니다. 구체화된 뷰는 여러 다운스트림 데이터 세트에 사용되는 변환, 집계 또는 사전 컴퓨팅 결과에 적합합니다.
보기 요청 시 평가되며 저장되지 않습니다. 중간 변환에 뷰를 사용하고 카탈로그에 게시할 필요가 없는 검사를 사용합니다.

스트리밍 테이블은 스트리밍 대상이기도 한 Unity 카탈로그 관리 테이블의 한 형태입니다. 스트리밍 테이블에는 하나 이상의 스트리밍 흐름(추가, AUTO CDC)이 기록되어 있을 수 있습니다. 스트리밍 흐름을 대상 스트리밍 테이블과 명시적으로 별도로 정의하거나 스트리밍 테이블 정의의 일부로 암시적으로 정의할 수 있습니다.

구체화된 뷰는 Unity 카탈로그 관리 테이블의 한 형태이기도 하며 일괄 처리 대상입니다. 구체화된 뷰에는 하나 이상의 구체화된 뷰 흐름이 기록될 수 있습니다. 구체화된 뷰는 항상 구체화된 뷰 정의의 일부로 흐름을 암시적으로 정의한다는 점에서 스트리밍 테이블과 다릅니다.

자세한 내용은 스트리밍 테이블구체화된 뷰를 참조하세요.

뷰, 구체화된 뷰 및 스트리밍 테이블을 사용하는 경우

파이프라인 쿼리를 구현할 때 사용 사례에 가장 적합한 데이터 세트 형식을 선택합니다.

뷰를 사용하여 다음을 수행할 수 있습니다.

  • 크거나 복잡한 쿼리를 관리하기 쉬운 쿼리로 분할합니다.
  • 기대치를 사용하여 중간 결과의 유효성을 검사합니다.
  • 유지할 필요가 없는 결과에 대한 스토리지 및 컴퓨팅 비용을 줄입니다. 테이블이 구체화되므로 추가 계산 및 스토리지 리소스가 필요합니다.

다음과 같은 경우 구체화된 뷰를 사용하는 것이 좋습니다.

  • 여러 다운스트림 쿼리는 테이블을 사용합니다. 구체화된 뷰는 결과를 캐시하므로 다운스트림 쿼리는 각 액세스에서 쿼리를 다시 계산하는 대신 미리 계산된 결과를 읽습니다.
  • 다른 파이프라인, 작업 또는 쿼리는 테이블을 사용합니다. 구체화된 뷰는 Unity 카탈로그 테이블로 구체화되므로 이를 정의하는 파이프라인 외부의 소비자는 쿼리할 수 있습니다. 뷰는 구체화되지 않으므로 동일한 파이프라인 내에서만 사용할 수 있습니다.
  • 개발 중에 쿼리 결과를 검사하려고 합니다. 구체화 뷰는 실제로 저장되며 파이프라인 외부에서도 쿼리할 수 있으므로, 개발 중에 계산 결과의 정확성을 검증할 수 있습니다. 유효성을 검사한 후 구체화할 필요가 없는 쿼리를 뷰로 변환합니다.
  • 쿼리가 집계 또는 조인을 수행하거나, 소스 데이터가 단순히 증가만 하는 것이 아니라 업데이트 및 삭제로 인해 변경될 수 있습니다. 구체화된 뷰는 원본 데이터의 현재 상태와 일치하는 결과를 유지하는 반면, 스트리밍 테이블은 추가 전용 원본용으로 설계되고 각 레코드를 한 번만 처리합니다.

다음과 같은 경우 스트리밍 테이블을 사용하는 것이 좋습니다.

  • 쿼리는 지속적으로 또는 증분적으로 증가하는 데이터 원본에 대해 정의됩니다.
  • 쿼리 결과는 증분 방식으로 계산되어야 합니다.
  • 파이프라인에는 높은 처리량과 짧은 대기 시간이 필요합니다.

Note

스트리밍 테이블은 항상 스트리밍 원본에 대해 정의됩니다. AUTO CDC ... INTO 스트리밍 원본을 사용하여 CDC 피드의 업데이트를 적용할 수도 있습니다. AUTO CDC API: 파이프라인을 사용하여 변경 데이터 캡처 간소화를 참조하세요.

Flows

흐름은 파이프라인의 기본 데이터 처리 개념이며 스트리밍 및 일괄 처리 의미 체계를 모두 지원합니다. 흐름은 원본에서 데이터를 읽고, 사용자 정의 처리 논리를 적용하고, 결과를 대상에 씁니다. 파이프라인은 Spark 구조적 스트리밍과 동일한 스트리밍 흐름 유형(추가, 업데이트, 완료)을 공유합니다. (현재 는 추가업데이트 흐름만 노출됩니다.) 자세한 내용은 구조적 스트리밍의 출력 모드를 참조하세요.

파이프라인은 다음과 같은 추가 흐름 형식도 제공합니다.

  • AUTO CDC 는 순서가 다른 CDC 이벤트를 처리하고 SCD Type 1 및 SCD Type 2를 모두 지원하는 Lakeflow 파이프라인의 고유한 스트리밍 흐름입니다. SDP에서는 자동 CDC를 사용할 수 없습니다.
  • 구체화된 뷰 는 가능한 한 새 데이터와 원본 테이블의 변경 내용만 처리하는 파이프라인의 일괄 처리 흐름입니다.

자세한 내용은 Lakeflow 파이프라인 흐름을 사용하여 증분 방식으로 데이터 로드 및 처리를 참조하세요.

Sinks

싱크는 파이프라인의 스트리밍 대상이며 델타 테이블, Apache Kafka 토픽, Azure EventHubs 토픽 및 사용자 지정 Python 데이터 원본을 지원합니다. 싱크에는 하나 이상의 스트리밍 흐름(추가, 업데이트)이 기록되어 있을 수 있습니다.

자세한 내용은 Lakeflow 파이프라인의 싱크를 참조하세요.

파이프라인

파이프라인은 개발 및 실행 단위이며 사용자가 정의한 흐름, 스트리밍 테이블, 구체화된 뷰 및 싱크에 대한 컨테이너입니다. 파이프라인 소스 코드에서 이러한 개체를 정의한 다음 파이프라인을 실행하여 파이프라인을 빌드합니다. 파이프라인이 실행되는 동안 정의된 개체의 종속성을 분석하고 실행 및 병렬화 순서를 자동으로 오케스트레이션합니다.

자세한 내용은 파이프라인이란?을 참조하세요.

또한 Lakeflow 파이프라인 외부에서 독립형 구체화된 뷰와 스트리밍 테이블을 정의할 수도 있으며, 이 경우 Azure Databricks가 파이프라인을 대신 관리합니다. 두 방법을 비교하려면 독립 실행형 파이프라인과 Lakeflow 파이프라인을 참조하세요.

파이프라인은 트리거 모드 또는 연속 모드로 실행되며, 이 모드는 사용 가능한 데이터를 새로고침하고 새 데이터가 도착할 때 테이블을 중지하거나 유지할지 제어합니다. 두 모드를 비교하려면 트리거 모드와 연속 파이프라인 모드를 참조하세요.

데이터 수집

파이프라인은 Azure Databricks에서 사용할 수 있는 모든 데이터 원본을 지원합니다. Databricks는 대부분의 수집 사용 사례에 스트리밍 테이블을 사용하는 것이 좋습니다. 클라우드 객체 스토리지의 파일에 대해 Auto Loader는 증분 및 멱등 로드를 제공합니다. 스트리밍 데이터의 경우 파이프라인은 Apache Kafka, Azure Event Hubs, Amazon Kinesis 및 Google Pub/Sub와 같은 메시지 버스에서 직접 수집할 수 있습니다. 파이프라인의 데이터 로드를 참조하세요.

데이터 품질

기대는 파이프라인을 통해 흐르는 데이터의 유효성을 검사하는 데이터 세트의 선택적 절입니다. 예상을 SQL 부울 제약 조건으로 정의하고 레코드가 실패할 때 발생하는 작업(경고, 레코드 삭제 또는 업데이트 실패)을 지정합니다. 파이프라인 기대를 사용하여 데이터 품질을 관리하기를 참조하세요.

델타 통합

파이프라인에서 만들고 관리하는 모든 테이블은 델타 테이블입니다. ACID 트랜잭션, 시간 이동 및 스키마 적용을 포함하여 Delta Lake와 동일한 보장을 제공합니다. 파이프라인은 추가 테이블 속성을 더하고 예측 최적화를 사용해 OPTIMIZEVACUUM 작업을 포함한 자동 유지 관리를 수행합니다. Azure Databricks에서 Delta Lake란?을 참조하세요.

추가 리소스