Native Lineage Support

Native Lineage Support (네이티브 계보 지원)

조직이 데이터 생태계를 관리(govern)하려 할 때, 데이터가 어디에서 오고 어디로 가는지 이해하는 데이터 계보(lineage)는 매우 중요해집니다. Apache Flink는 스트리밍 데이터 레이크에서 데이터 수집과 ETL에 널리 사용되므로, 다음을 포함한(이에 국한되지 않는) 시나리오를 위한 종단 간 계보 솔루션이 필요합니다.

출처: 문서

본문

  • Data Quality Assurance (데이터 품질 보증): 데이터 파이프라인 내에서 데이터 오류를 그 기원까지 추적하여 데이터 불일치를 식별하고 바로잡습니다.
  • Data Governance (데이터 거버넌스): 데이터의 출처와 변환을 문서화하여 명확한 데이터 소유권과 책임을 확립합니다.
  • Regulatory Compliance (규제 준수): 수명 주기 전반에 걸쳐 데이터 흐름과 변환을 추적하여 데이터 프라이버시 및 규제 준수를 보장합니다.
  • Data Optimization (데이터 최적화): 중복 데이터 처리 단계를 식별하고 데이터 흐름을 최적화하여 효율성을 높입니다.

Apache Flink는 개발자가 계보 메타데이터를 OpenLineage 같은 외부 계보 시스템에 통합할 수 있도록, 내부 계보 데이터 모델과 Job Status Listener를 제공하여 네이티브 계보 지원을 제공합니다. Flink 런타임에서 job이 생성되면 JobCreatedEvent에 Lineage Graph 메타데이터가 포함되며, 이는 Job Status Listener로 전송됩니다.

Lineage 데이터 모델 (Lineage Data Model)

Flink 네이티브 계보 인터페이스는 두 계층으로 정의됩니다. 첫 번째 계층은 모든 Flink job과 커넥터를 위한 일반(generic) 인터페이스이고, 두 번째 계층은 Table과 DataStream 각각에 대한 확장 인터페이스를 정의합니다. 인터페이스와 클래스 간의 관계는 아래 다이어그램에 정의되어 있습니다.

기본적으로 Flink Table 환경에서는 Table 관련 계보 인터페이스 또는 클래스가 사용되므로, Flink 사용자는 이러한 인터페이스를 건드릴 필요가 없습니다. Flink 커뮤니티는 Kafka, JDBC, Cassandra, Hive와 같은 공통 커넥터를 점차 지원할 예정입니다. 커스텀 커넥터를 정의했다면 LineageVertexProvider 인터페이스의 커스텀 source/sink 구현이 필요합니다. LineageVertex 안에는 Flink source/sink의 메타데이터로 Lineage Dataset 목록이 정의됩니다.

@PublicEvolving
public interface LineageVertexProvider {
  LineageVertex getLineageVertex();
}

인터페이스에 대한 자세한 내용은 FLIP-314를 참고하세요.

더 알아보기 (Learn more)