파이프라인에 대한 단위 테스트

Important

이 기능은 베타 버전으로 제공됩니다.

Databricks의 Python 단위 테스트에 대한 일반적인 내용은 Python 단위 테스트를 참조하세요.

Lakeflow 파이프라인은 웹 기반 Lakeflow 파이프라인 편집기에서 Python 단위 테스트 작성을 지원합니다. 이렇게 하면 모의 데이터를 사용하여 Python 또는 SQL 변환 논리의 유효성을 검사할 수 있습니다. 파이프라인 테스트 프레임워크를 사용하면 에지 사례를 테스트하고, 독점 파이프라인 API(자동 CDC, 스트리밍 테이블, 예상, 추가 흐름)의 유효성을 검사하고, 지원되는 테이블 식별자 작업에 모의 입력을 사용하여 반복할 수 있습니다. 테스트를 실행하기 전에 격리 제한을 검토합니다.

  • 격리된 테스트 실행: 프레임워크는 파이프라인의 기본 카탈로그에서 테이블 작업을 임시 테스트 스키마로 리디렉션하는 SparkSession을 제공하므로 프로덕션 테이블에 영향을 주지 않고 입력 데이터를 모의하고 테스트 출력을 작성할 수 있습니다. 격리는 이름으로 테이블을 참조하는 작업에 적용됩니다. 제한 사항을 참조하세요.
  • 유연한 테스트 범위: 테스트 SparkSession을 사용하여 파이프라인의 컴퓨팅에서 파이프라인의 하위 집합(개별 테이블, 종속 테이블 체인 또는 전체 파이프라인)을 실행합니다.
  • 결과 유효성 검사: 표준 pytest 어설션을 사용하여 테스트에서 생성된 격리된 출력 테이블의 결과를 확인합니다.

단위 테스트를 사용하는 경우

일반적인 사용 사례는 다음과 같습니다.

  • 새 변환 논리 유효성 검사: 프로덕션 데이터에 대해 실행하기 전에 변환에서 예상 스키마, 행 수, 집계 및 비즈니스 논리를 생성하는지 테스트합니다.
  • 자동 CDC 사양 테스트: 자동 CDC 흐름 정의가 모의 데이터를 사용하여 변경 이벤트를 올바르게 처리하고 삽입, 업데이트, 삭제 및 SCD(느린 변경 차원) 형식을 처리하는지 확인합니다.
  • 기대치 및 데이터 품질 규칙 테스트: 기대치가 실패해야 할 때 실패하고 데이터가 유효할 때 통과하는지 확인합니다.
  • 종속 테이블 간 테스트: 변환 체인(예: 브론즈, 실버 및 골드)을 테스트하여 파이프라인 그래프를 통해 데이터가 올바르게 흐르는지 확인합니다.

Requirements

  • 파이프라인 Owner 권한 그리고 파이프라인의 기본 카탈로그에 대한 USE CATALOGCREATE SCHEMA 권한 테스트가 실행되는 임시 테스트 스키마를 만들려면 프레임워크에 이러한 권한이 필요합니다.

    파이프라인 권한을 확인하거나 설정하려면 파이프라인을 열고 공유를 클릭합니다. 파이프라인은 Owner(IS OWNER)이어야 합니다. CAN RUNCAN MANAGE만으로는 테스트를 실행할 수 없습니다. 파이프라인 권한 구성을 참조하세요.

    카탈로그 권한을 확인하거나 설정하려면 카탈로그 탐색기에서 카탈로그를 열고 사용 권한 탭을 선택하고 가지고 있는지 USE CATALOG 확인합니다 CREATE SCHEMA. 카탈로그 소유자, 메타스토어 관리자 또는 권한이 있는 MANAGE 사용자는 SQL을 포함하여 부여할 수 있습니다.

    GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;
    

    자세한 내용은 Unity 카탈로그 권한 참조를 참조하세요.

  • 파이프라인은 트리거(비연속) 모드로 구성해야 합니다.

  • 파이프라인은 Databricks 런타임 18.1 이상에서 실행되어야 합니다. 초기 런타임에는 단위 테스트 모듈이 포함되어 있지 않습니다. 업데이트가 어떤 런타임 버전에서 실행되었는지 확인하려면 파이프라인 이벤트 로그를 쿼리하세요. 런타임 정보를 참조하세요.

  • Spark Connect는 지원되지 않습니다.

메모

테스트 격리는 테이블 이름으로 테이블을 참조하는 작업을 포함합니다. 격리를 우회하는 작업은 테스트 코드와 선택한 출력에서 실행되는 파이프라인 코드에서 전이적 종속성을 포함하여 발생할 수 있습니다. 안전해 보이는 테스트 파일은 프로덕션 데이터에서 작동하는 경로 또는 커넥터를 통해 읽거나 쓰는 파이프라인 흐름을 계속 실행할 수 있습니다. 테스트가 프로덕션 데이터 또는 메타데이터에 영향을 주지 않도록 하려면 다음 규칙을 따릅니다.

  • 모든 테이블을 이름(catalog.schema.table)으로 참조하고 모든 입력을 이름으로 모의합니다. 경로(/Volumes/..., , dbfs:/...,s3://...abfss://...)로 읽거나 쓰지 않으며 Kafka 또는 자동 로더와 같은 커넥터에서 읽지 마세요. 이러한 격리를 우회하고 실제 프로덕션 시스템에서 작동합니다.
  • 거버넌스 또는 소유권 문(예: GRANT, REVOKE,ALTER ... OWNER TOSET/UNSET TAGS 또는 CREATE/DROP POLICY.)을 실행하지 마세요. 이것들은 실제 운영 보안 개체를 대상으로 실행됩니다.
  • 카탈로그 또는 스키마(CREATE CATALOG, CREATE SCHEMA)를 만들지 마세요. 이는 실제 Unity Catalog 메타스토어에 연결됩니다.
  • 그래프에 경로 기반 입력, 커넥터, 명령적 쓰기 또는 기타 외부 부작용이 포함된 경우 전체 파이프라인을 실행하지 마세요. 종속성이 지원되는 카탈로그 테이블 작업을 사용하고 모의 입력으로 대체된 출력만 선택합니다.

자세한 내용은 제한 사항을 참조하세요.

Limitations

Warning

일부 작업은 테스트 격리를 무시하고 실제 프로덕션 데이터 또는 메타데이터에 대해 작동할 수 있습니다. 테스트를 실행하기 전에 다음 제한 사항을 검토합니다.

테스트 격리는 테이블 이름에만 해당합니다.

  • 경로 또는 커넥터를 통해 읽거나 쓰지 마세요. 격리는 이름(예 spark.read.table("catalog.schema.table") : 또는 df.write.saveAsTable("catalog.schema.table"))으로 테이블을 참조하는 작업만 리디렉션합니다. 경로 또는 커넥터를 통해 처리되는 작업은 격리를 우회하고 실제 프로덕션 시스템에서 직접 작동합니다.

    • 경로(예: df.write.save("/Volumes/...")dbfs:/ 경로 또는 클라우드 또는 외부 위치 경로(예: s3://... 또는abfss://...)를 기준으로 작성하면 실제 프로덕션 스토리지에 쓰고 프로덕션 데이터를 덮어쓸 수 있습니다.
    • 경로별 읽기 (예: spark.read.load(path) 또는 spark.read.format("delta").load(path))는 모의 데이터 대신 실제 프로덕션 데이터를 반환합니다.
    • 커넥터에서 읽기는 실제 프로덕션 소스에 연결됩니다. 여기에는 Kafka (실제 브로커에서 읽기) 및 자동 로더 (cloudFiles실제 클라우드 스토리지 경로에서 읽음)가 포함됩니다. 둘 다 모의 데이터로 리디렉션되지 않습니다.
  • 파이프라인 단위 테스트에서 event_log()테이블 값 함수를 사용하지 마세요. 테스트 모드에서는 event_log() 테스트 실행의 이벤트 로그로 리디렉션되지 않습니다. 운영 환경의 이벤트 로그 또는 이전에 등록된 이벤트 로그를 반환할 수 있으므로, 이를 대상으로 한 어설션이 운영 데이터를 읽게 될 수 있습니다. 대신 run에서 반환된 event_log_table_name를 사용하고 test_spark를 통해 이를 조회하세요. event_log_table_nameNone (예: 이벤트 로그 테이블 이름을 확인할 수 없는 경우) 쿼리하기 전에 확인합니다.

    status = test_pipeline.run(test_spark, set(["catalog.schema.table"]))
    assert status.event_log_table_name is not None
    events = test_spark.table(status.event_log_table_name)
    

    실패한 업데이트를 진단하는 것이 목표인 경우 이벤트 로그를 읽기 전에 어설션 status.is_success 하지 마세요. 이벤트 로그는 업데이트가 실패한 이유를 이해하기 위해 검사하는 경우가 많습니다.

거버넌스 및 DDL 작업

  • 카탈로그, 스키마, 권한, 소유권, 태그 및 정책 변경은 지원되지 않습니다. 여기에는 CREATE/DROP/ALTER CATALOG,CREATE/DROP/ALTER SCHEMA(포함SET MANAGED LOCATION),GRANT/REVOKE , ALTER ... OWNER TO및 . SET/UNSET TAGSCREATE/DROP POLICY 실행 test_spark 된 일부 SQL 양식은 심층 방어로 거부됩니다. 다른 양식 또는 직접 API를 통해 호출된 동일한 작업은 실제 프로덕션 개체에 도달할 수 있습니다. 이러한 가드를 격리 경계로 사용하지 마세요. 이러한 문을 테스트 코드와 선택한 출력에서 실행된 파이프라인 코드에서 제외합니다.

운영 제한 사항

  • 동시 실행은 지원되지 않습니다. 테스트 및 파이프라인 업데이트를 동시에 실행하는 것은 지원되지 않으며 시스템에서 이를 방지하지 않습니다. 둘 사이에는 조정이 없으므로 동시에 실행하면 리소스에 대해 경합할 수 있으며 프로덕션 업데이트의 성능이 심각하게 저하되거나 테스트가 시작되지 않습니다. 파이프라인이 업데이트를 실행하는 동안 테스트를 시작하지 마세요(또는 테스트가 실행되는 동안 업데이트를 시작). 테스트를 실행하기 전에 진행 중인 업데이트가 완료되기를 기다립니다.
  • 비정상적인 종료 후 임시 스키마: 각 테스트 실행은 파이프라인의 기본 카탈로그에 임시 스키마(명명됨 redirecting_<id>)를 만들고 실행이 완료되면 자동으로 삭제합니다. 실행이 비정상적으로 종료되는 경우(예: 컴퓨팅이 중간 실행에서 손실됨) 임시 스키마는 런의 모의 테이블과 출력 테이블을 유지하여 남겨둘 수 있습니다. 프로덕션 데이터에는 영향을 주지 않습니다. 스토리지를 회수하려면 파이프라인의 기본 카탈로그에서 이름이 시작되는 redirecting_ 남은 스키마를 수동으로 삭제합니다.
  • 테스트 실행은 컴퓨팅을 사용합니다. 테스트 실행은 파이프라인의 컴퓨팅에서 실행되며 일반 파이프라인 업데이트로 청구됩니다. 테스트 실행에 대한 별도의 계량은 없습니다.
  • 전체 새로 고침은 지원되지 않습니다. 선택적 새로 고침만 사용할 수 있습니다. test_pipeline.run() 선택한 출력(또는 선택 항목을 전달하지 않을 때 모든 출력)을 새로 고칩니다. 전체 새로 고침 및 전체 새로 고침 선택은 구현되지 않습니다.

작성 및 충실도 제한 사항

  • 편집기 전용 실행: 웹 기반 Lakeflow 파이프라인 편집기에서 테스트를 실행해야 합니다.
  • Python 테스트만: 테스트는 Python 작성해야 합니다. SQL 파이프라인을 테스트할 수 있지만 테스트 자체는 Python 작성해야 합니다.
  • 거버넌스 충실도: 모의 데이터는 대체되는 프로덕션 테이블에 정의된 행 필터 또는 열 마스크를 상속하지 않습니다. 테스트 결과는 모의 입력을 제공하는 것과 정확하게 반영하며, 관리되는 프로덕션 데이터에서 동일한 쿼리가 동작하는 방식과 다를 수 있습니다.

1단계: 파이프라인 설정 업데이트

파이프라인을 트리거 모드로 설정하세요.

  1. UI에서 파이프라인을 열고 설정을 클릭하세요.
  2. 파이프라인 모드트리거됨으로 설정합니다(연속 모드를 사용하지 않음).

또는 파이프라인 설정 JSON을 직접 편집합니다.

"continuous": false

2단계: 테스트 파일 만들기

Lakeflow 파이프라인 편집기에서 (추가) 단추를 클릭하고 +테스트를 선택합니다. 이렇게 하면 파이프라인 소스 코드에 tests 포함되지 않은 테스트 파일(및 폴더가 아직 없는 경우)이 만들어집니다. tests 폴더를 직접 만들 필요가 없습니다.

테스트 옵션을 보여 주는 파이프라인 자산 메뉴를 추가하여 pytest 파일을 만듭니다.

3단계: 테스트 생성

Genie Code는 테스트 스캐폴딩을 생성할 수 있습니다.

  • 테스트 파일 내에서 테스트 생성 단추를 클릭합니다.

    테스트 생성 단추가 있는 빈 테스트 파일입니다.

  • 대신 Genie Code 에이전트 모드에서 /tests를 사용할 수 있습니다.

    TestPipeline 기반 단위 테스트를 사용하여 Genie Code로 채워진 테스트 파일입니다.

Genie Code를 사용해 상용구 코드를 생성한 다음, 예외 사례에 맞게 맞춤 설정합니다.

또는 테스트 코드를 직접 작성할 수 있습니다. 각 테스트 파일의 맨 위에 다음 가져오기를 추가합니다.

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark

test_pipeline = TestPipeline.active()

4단계: 테스트 실행

Lakeflow 파이프라인 편집기에서 테스트를 실행합니다.

  • 테스트 함수 옆의 여백에서 재생 아이콘 (재생) 단추를 클릭하여 개별 테스트를 실행합니다.
  • 테스트 파일 맨 위에 있는 파일에서 테스트 실행을 클릭하여 해당 파일의 모든 테스트를 실행합니다.

테스트 결과(성공 또는 실패)가 편집기 아래쪽 패널에 표시됩니다. 어설션 오류를 검토하여 오류를 디버그합니다.

API 테스트

API 설명
TestPipeline.active() TestPipeline 현재 Lakeflow 파이프라인 편집기에서 편집 중인 파이프라인의 개체를 반환합니다. 이 개체는 소스 코드, 구성, 기본 카탈로그/스키마 등을 포함한 파이프라인에 대한 참조입니다.
test_pipeline.run(test_spark, set([table_names])) 테이블 이름이 지정된 경우 선택적 새로 고침을 수행하여 파이프라인의 업데이트를 동기적으로 실행합니다. 파이프라인 실행이 성공하거나 예외로 종료된 후 반환합니다.
test_spark 고정 장치 테이블을 이름으로 참조하는 테이블 읽기/쓰기 작업(예: spark.read.table("catalog.schema.table") 또는 df.write.saveAsTable("catalog.schema.table"))을 임시 테스트 스키마로 자동 리디렉션하는 카탈로그 테이블 리디렉션이 적용된 테스트 SparkSession을 생성합니다. 리디렉션은 이름 기반 테이블 작업에만 적용됩니다. 실제 시스템에서 직접 작동하는 경로 또는 커넥터를 통해 주소가 지정된 읽기 또는 쓰기 는 다루지 않습니다 . 제한 사항을 참조하세요.

모의 데이터 만들기

SQL 또는 createDataFrame다음 중 하나를 사용하여 입력 데이터를 모의로 만들 수 있습니다.

# Option 1: Using SQL
test_spark.sql("""
    CREATE TABLE catalog.schema.table_name AS
    SELECT * FROM VALUES
        (1, 'value1'),
        (2, 'value2')
    AS t(id, name)
""")

# Option 2: Using createDataFrame
df = test_spark.createDataFrame(
    [(1, 'value1'), (2, 'value2')],
    schema=["id", "name"]
)
df.write.saveAsTable("catalog.schema.table_name")

더 많은 양의 실제 가상 데이터를 생성하려면 Faker 라이브러리를 사용할 수 있습니다. 먼저 파이프라인에서 %pip install faker를 실행한 다음, Faker 기반 UDF로 DataFrame을 생성하세요.

# Option 3: Using Faker for synthetic data
from pyspark.sql import functions as F
from faker import Faker

fake = Faker()
fake_firstname = F.udf(fake.first_name)
fake_lastname = F.udf(fake.last_name)
fake_email = F.udf(fake.ascii_company_email)

df = (
    test_spark.range(0, 100)
    .withColumn("firstname", fake_firstname())
    .withColumn("lastname", fake_lastname())
    .withColumn("email", fake_email())
)
df.write.saveAsTable("catalog.schema.table_name")

파이프라인 또는 특정 테이블 실행

# Run specific tables
test_pipeline.run(test_spark, set(["catalog.schema.table1", "catalog.schema.table2"]))

# Run all tables in the pipeline
test_pipeline.run(test_spark)

예제

예제 1: 행 수, 스키마 및 null 처리를 사용하여 집계 테스트

목표: 사용자 집계의 유효성을 검사하여 유형별로 사용자를 올바르게 계산하고, null 전자 메일을 처리하고, 예상된 스키마를 생성합니다.

파이프라인 변환:

이러한 변환은 간단한 두 테이블 파이프라인 users 을 만듭니다. 즉, counts 사용자 데이터를 선택하고 유형별로 사용자를 그룹화하고 총 사용자 및 유효한 전자 메일 수를 계산합니다.

from pyspark import pipelines as dp
from pyspark.sql.functions import col, count, count_if

@dp.table
def users():
    return (
        spark.read.table("catalog.schema.wanderbricks_users")
        .select("user_id", "email", "name", "user_type")
    )

@dp.table
def counts():
    return (
        spark.read.table("catalog.schema.users")
        .withColumn("valid_email", col("email").isNotNull())
        .groupBy("user_type")
        .agg(
            count("user_id").alias("total_count"),
            count_if("valid_email").alias("count_valid_emails")
        )
    )

테스트:

이러한 테스트는 의도적인 null을 사용하여 모의 사용자 데이터를 만들고 파이프라인을 격리된 상태로 실행하여 행 수, 스키마 구조, null 처리 및 집계 논리의 유효성을 검사합니다.

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
from pyspark.testing import assertDataFrameEqual

test_pipeline = TestPipeline.active()

# Mock data fixture
def mock_users(session):
    session.sql("""
        CREATE TABLE catalog.schema.wanderbricks_users AS
        SELECT * FROM VALUES
            (1, 'alice@example.com', 'Alice', 'admin'),
            (2, NULL, 'Bob', 'user'),
            (3, 'charlie@example.com', 'Charlie', 'user'),
            (4, NULL, 'Dana', 'admin')
        AS t(user_id, email, name, user_type)
    """)

# Test 1: Row count
def test_users_row_count(test_spark):
    mock_users(test_spark)
    test_pipeline.run(test_spark, set(["catalog.schema.users"]))
    result = test_spark.table("catalog.schema.users")
    assert result.count() == 4

# Test 2: Schema validation
def test_users_schema(test_spark):
    mock_users(test_spark)
    test_pipeline.run(test_spark, set(["catalog.schema.users"]))
    result = test_spark.table("catalog.schema.users")
    expected_fields = {"user_id", "email", "name", "user_type"}
    actual_fields = set(f.name for f in result.schema.fields)
    assert expected_fields == actual_fields

# Test 3: Null handling
def test_users_null_handling(test_spark):
    mock_users(test_spark)
    test_pipeline.run(test_spark, set(["catalog.schema.users"]))
    result = test_spark.table("catalog.schema.users")
    null_emails = result.filter("email IS NULL").count()
    assert null_emails == 2

# Test 4: Aggregation
def test_counts(test_spark):
    mock_users(test_spark)
    # Run both tables since counts depends on users
    test_pipeline.run(test_spark, set(["catalog.schema.users", "catalog.schema.counts"]))
    result = test_spark.table("catalog.schema.counts")
    # Check counts for each user_type
    admin_row = result.filter("user_type = 'admin'").collect()[0]
    user_row = result.filter("user_type = 'user'").collect()[0]
    assert admin_row["total_count"] == 2
    assert admin_row["count_valid_emails"] == 1
    assert user_row["total_count"] == 2
    assert user_row["count_valid_emails"] == 1

# Test 5: Full DataFrame comparison with assertDataFrameEqual
def test_counts_full_dataframe(test_spark):
    mock_users(test_spark)
    test_pipeline.run(test_spark, set(["catalog.schema.users", "catalog.schema.counts"]))
    result = test_spark.table("catalog.schema.counts")
    expected = test_spark.createDataFrame(
        [("admin", 2, 1), ("user", 2, 1)],
        schema=["user_type", "total_count", "count_valid_emails"]
    )
    assertDataFrameEqual(result, expected)

예제 2: 자동 CDC 테스트

목표: 자동 CDC가 삽입 및 업데이트를 사용하여 변경 피드를 올바르게 처리하는지 확인합니다.

파이프라인 변환:

이 변환은 변경 피드에서 자동 CDC를 설정하여 스트리밍 변경 내용을 읽고 대상 테이블에 SCD 유형 1로 적용합니다(최신 버전만 유지).

from pyspark import pipelines as dp
from pyspark.sql.functions import col

@dp.view
def users():
    return spark.readStream.table("catalog.schema.change_feed")

dp.create_streaming_table("target_autocdc")
dp.create_auto_cdc_flow(
    target="target_autocdc",
    source="users",
    keys=["userId"],
    sequence_by=col("ts"),
    stored_as_scd_type=1
)

테스트:

첫 번째 테스트는 동일한 userId 레코드(업데이트 시뮬레이트)에 대해 여러 레코드가 있는 모의 변경 피드를 만들고 최신 레코드만 대상에 유지되도록 확인합니다. 두 번째 테스트는 파이프라인을 실행하고, 변경 피드에 더 많은 이벤트를 추가하고, 파이프라인을 다시 실행하여 지연 도착 및 잘못된 순서 이벤트를 시뮬레이션합니다.

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark

test_pipeline = TestPipeline.active()

# Test 1: Standard inserts and updates
def test_auto_cdc_flow(test_spark):
    # Create a mock change feed table
    test_spark.sql("""
        CREATE TABLE catalog.schema.change_feed AS
        SELECT * FROM VALUES
            (1, 'Alice', 1000),
            (2, 'Bob', 1001),
            (1, 'Alice Updated', 1002)
        AS t(userId, name, ts)
    """)
    # Run the pipeline
    test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))
    # Read the output
    result = test_spark.table("catalog.schema.target_autocdc")
    # Verify two users exist
    user_ids = set(row["userId"] for row in result.collect())
    assert user_ids == {1, 2}
    # Verify latest record for userId=1 has ts=1002
    latest_user1 = result.filter("userId = 1").collect()[0]
    assert latest_user1["ts"] == 1002
    assert latest_user1["name"] == "Alice Updated"
    # Verify userId=2 has ts=1001
    user2 = result.filter("userId = 2").collect()[0]
    assert user2["ts"] == 1001

# Test 2: Late-arriving and out-of-order events
def test_auto_cdc_late_arriving(test_spark):
    # First batch of change events
    test_spark.sql("""
        CREATE TABLE catalog.schema.change_feed AS
        SELECT * FROM VALUES
            (1, 'Alice', 1000),
            (2, 'Bob', 1001)
        AS t(userId, name, ts)
    """)
    # Run the pipeline with the initial batch
    test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))

    # Append late-arriving events to the change feed:
    # - A newer event for userId=1 (ts=1003) that arrived after the first run
    # - A stale event for userId=2 (ts=999) with a timestamp older than what is already applied
    test_spark.sql("""
        INSERT INTO catalog.schema.change_feed VALUES
            (1, 'Alice Updated', 1003),
            (2, 'Bob (stale)', 999)
    """)
    # Re-run the pipeline. sequence_by=ts ensures stale events do not overwrite newer state.
    test_pipeline.run(test_spark, set(["catalog.schema.target_autocdc"]))

    result = test_spark.table("catalog.schema.target_autocdc")
    # userId=1 should reflect the newer late-arriving event
    alice = result.filter("userId = 1").collect()[0]
    assert alice["ts"] == 1003
    assert alice["name"] == "Alice Updated"
    # userId=2 should be unchanged: the stale event with an older ts is ignored
    bob = result.filter("userId = 2").collect()[0]
    assert bob["ts"] == 1001
    assert bob["name"] == "Bob"

예제 3: 스냅샷에서 자동 CDC 테스트

목표: CDC가 삽입, 업데이트 및 삭제를 포함한 스냅샷 변경 내용을 올바르게 처리하는지 확인합니다.

파이프라인 변환:

이 변환은 스냅샷 테이블에서 읽고 시간이 지남에 따른 변경 사항을 SCD 유형 2(전체 이력 유지)로 추적하는 스냅샷 기반 Auto CDC를 구성합니다.

from pyspark import pipelines as dp

@dp.view(name="source")
def source():
    return spark.read.table("catalog.schema.snapshot")

dp.create_streaming_table("catalog.schema.target")
dp.create_auto_cdc_from_snapshot_flow(
    target="target",
    source="source",
    keys=["userId"],
    stored_as_scd_type=2
)

테스트:

이 테스트는 초기 스냅샷을 만들고, 파이프라인을 실행한 다음, CDC가 모든 변경 내용을 캡처하는지 확인하기 위해 새 데이터를 잘리고 삽입하여 스냅샷 업데이트를 시뮬레이션합니다.

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark

test_pipeline = TestPipeline.active()

def test_auto_cdc_from_snapshot_flow(test_spark):
    # Create initial snapshot
    test_spark.sql("""
        CREATE TABLE catalog.schema.snapshot AS
        SELECT * FROM VALUES
            (1, 'Alice', '2024-01-01'),
            (2, 'Bob', '2024-01-02')
        AS t(userId, name, created_at)
    """)
    # Run the pipeline
    test_pipeline.run(test_spark, set(["catalog.schema.target"]))
    # Simulate a new snapshot by truncating and inserting updated data
    test_spark.sql("TRUNCATE TABLE catalog.schema.snapshot")
    test_spark.sql("INSERT INTO catalog.schema.snapshot VALUES (2, 'Bob', '2024-01-03')")
    test_pipeline.run(test_spark, set(["catalog.schema.target"]))
    # Verify SCD Type 2: should have 3 rows (original Alice, original Bob, updated Bob)
    result = test_spark.table("catalog.schema.target")
    assert result.count() == 3
    user_ids = [row["userId"] for row in result.collect()]
    assert set(user_ids) == {1, 2}

예제 4: 조인 및 예상 테스트

목표: 조인이 올바르게 작동하고 기대가 잘못된 데이터를 필터링하는지 확인합니다.

파이프라인 변환:

이 변환은 숙소 이미지를 편의시설과 조인한 후 2024년 1월 이전에 업로드된 이미지를 필터링하는 조건을 적용합니다.

from pyspark import pipelines as dp

@dp.table
@dp.expect_or_drop("uploaded after Jan 2024", "uploaded_at > '2024-01-01'")
def property_images_amenities_join():
    return (
        spark.read.table("catalog.schema.property_images")
        .join(
            spark.read.table("catalog.schema.property_amenities"),
            on="property_id",
            how="inner"
        )
    )

테스트:

이러한 테스트는 조인이 올바른 수의 행을 생성하는지, 그리고 기대값이 잘못된 업로드 날짜가 있는 레코드를 성공적으로 걸러내는지를 검증합니다.

import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark

test_pipeline = TestPipeline.active()

# Mock property datasets
def mock_properties(session):
    session.sql("""
        CREATE TABLE catalog.schema.property_images AS
        SELECT * FROM VALUES
            (101, 'img1.jpg', '2024-02-01'),
            (102, 'img2.jpg', '2024-01-15'),
            (103, 'img3.jpg', '2024-12-20')
        AS t(property_id, image_url, uploaded_at)
    """)
    session.sql("""
        CREATE TABLE catalog.schema.property_amenities AS
        SELECT * FROM VALUES
            (101, 'wifi'),
            (102, 'pool'),
            (103, 'parking')
        AS t(property_id, amenity)
    """)

# Test 1: Join
def test_property_join(test_spark):
    mock_properties(test_spark)
    test_pipeline.run(test_spark, set(["catalog.schema.property_images_amenities_join"]))
    result = test_spark.table("catalog.schema.property_images_amenities_join")
    # Should have 3 rows after join
    assert result.count() == 3
    # Check all property_ids are present
    property_ids = set(row["property_id"] for row in result.collect())
    assert property_ids == {101, 102, 103}

# Test 2: Expectation
def test_property_expectation(test_spark):
    mock_properties(test_spark)
    # Add a row with uploaded_at before Jan 2024
    test_spark.sql("""
        INSERT INTO catalog.schema.property_images VALUES (104, 'img4.jpg', '2023-12-31')
    """)
    # Add a matching row in the amenities table for the join
    test_spark.sql("""
        INSERT INTO catalog.schema.property_amenities VALUES (104, 'gym')
    """)
    test_pipeline.run(test_spark, set(["catalog.schema.property_images_amenities_join"]))
    result = test_spark.table("catalog.schema.property_images_amenities_join")
    # Only property_ids with uploaded_at > '2024-01-01' should be present
    valid_ids = set(row["property_id"] for row in result.collect())
    assert 104 not in valid_ids
    assert valid_ids == {101, 102, 103}