パイプラインの単体テスト

Important

この機能は ベータ版です

Databricks での単体テストPythonに関する一般的な情報については、単体テストPython参照してください。

Lakeflow パイプラインでは、Web ベースの Lakeflow パイプライン エディター Python単体テストの記述がサポートされています。 これにより、モック データを使用してPythonまたは SQL 変換ロジックを検証できます。 パイプライン テスト フレームワークを使用すると、エッジ ケースをテストし、独自のパイプライン API (自動 CDC、ストリーミング テーブル、期待、追加フロー) を検証し、サポートされているテーブル識別子操作のモック入力を使用して反復処理を行うことができます。 テストを実行する前に、分離の制限事項を確認してください。

  • 分離されたテストの実行: フレームワークは、パイプラインの既定のカタログ内の一時的なテスト スキーマにテーブル操作をリダイレクトする SparkSession を提供するため、運用テーブルに影響を与えずに入力データをモックし、テスト出力を書き込むことができます。 分離は、名前でテーブルを参照する操作に適用されます。 制限事項を参照してください。
  • 柔軟なテスト スコープ: テスト SparkSession を使用して、パイプラインのサブセット (個々のテーブル、依存テーブルのチェーン、またはパイプライン全体) をパイプラインのコンピューティングで実行します。
  • 結果の検証: 標準の pytest アサーションを使用して、テストで作成された分離出力テーブルの結果を確認します。

単体テストを使用する場合

一般的なユース ケースは次のとおりです。

  • 新しい変換ロジックの検証: 運用データに対して実行する前に、変換によって予期されるスキーマ、行数、集計、およびビジネス ロジックが生成されることをテストします。
  • 自動 CDC 仕様のテスト: 自動 CDC フロー定義が、モック データを使用して変更イベントを正しく処理し、挿入、更新、削除、SCD (緩やかに変化するディメンション) の種類を処理していることを検証します。
  • 期待条件とデータ品質ルールのテスト: 期待条件が失敗すべきときに失敗し、データが有効な場合に成功することを確認します。
  • 依存テーブル間のテスト: 変換のチェーン (ブロンズ、シルバー、ゴールドなど) をテストして、パイプライン グラフを通じてデータが正しく流れるかどうかを検証します。

Requirements

  • パイプラインOwner権限に加えて、パイプラインの既定のカタログに対するUSE CATALOGおよびCREATE SCHEMA権限 フレームワークには、テストを実行する一時的なテスト スキーマを作成するために、これらの特権が必要です。

    パイプラインのアクセス許可を確認または設定するには、パイプラインを開き、[ 共有] をクリックします。 パイプライン Owner (IS OWNER) である必要があります。 CAN RUNCAN MANAGE では、テストを実行するには不十分です。 「パイプラインのアクセス許可を構成する」を参照してください。

    カタログ権限を確認または設定するには、カタログ エクスプローラーでカタログを開き、[ アクセス許可 ] タブを選択し、 USE CATALOGCREATE SCHEMAがあることを確認します。 カタログ所有者、メタストア管理者、または MANAGE 特権を持つユーザーは、SQL を含め、それらを付与できます。

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

    詳細については、「 Unity カタログの権限リファレンス」を参照してください

  • パイプラインは、トリガー (非連続) モードで構成する必要があります。

  • パイプラインは プレビュー チャネル上にある必要があります。 単体テストはベータ版であり、プレビューでのみ利用できます。

  • Spark Connect はサポートされていません。

Note

テスト分離では、 テーブルを名前で参照するテーブル操作について説明します。 分離をバイパスする操作は、テスト コードと、選択した出力によって実行されるパイプライン コード (推移的な依存関係を含む) の両方で発生する可能性があります。 安全に見えるテスト ファイルでは、運用データに対して動作するパスまたはコネクタによって読み取りまたは書き込みを行うパイプライン フローを実行できます。 テストが運用データまたはメタデータに影響しないようにするには、次の規則に従います。

  • すべてのテーブルを名前 (catalog.schema.table) で参照し、すべての入力を名前でモックします。 パス (/Volumes/...dbfs:/...s3://...abfss://...) による読み取りまたは書き込み、Kafka や自動ローダーなどのコネクタからの読み取りは行わないでください。 これらは分離をバイパスし、実際の運用システムで動作します。
  • GRANTREVOKEALTER ... OWNER TOSET/UNSET TAGSCREATE/DROP POLICYなどのガバナンスステートメントや所有権ステートメントは実行しないでください。 これらは、実際の運用セキュリティ保護可能なリソースに対して実行されます。
  • カタログまたはスキーマ (CREATE CATALOGCREATE SCHEMA) は作成しないでください。 これらは、実際の Unity カタログ メタストアに到達します。
  • そのグラフにパスベースの入力、コネクタ、命令型書き込み、またはその他の外部副作用が含まれている場合は、パイプライン全体を実行しないでください。 サポート対象のカタログテーブル操作を使用し、モック入力に置き換えられている依存関係を持つ出力のみを選択します。

詳細については、「制限事項」を参照してください。

制限事項

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() はテスト実行のイベント ログにリダイレクトされません。 運用環境または以前に登録されたイベント ログを返すことができるので、それに対するアサーションが運用データを読み取る可能性があります。 代わりに、実行によって返される 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 CATALOGCREATE/DROP/ALTER SCHEMA ( SET MANAGED LOCATIONを含む)、 GRANT/REVOKEALTER ... OWNER TOSET/UNSET TAGS、および CREATE/DROP POLICYが含まれます。 test_sparkによって実行される一部の SQL フォームは、多層防御として拒否されます。他のフォーム、または直接 API を介して呼び出された同じ操作は、実際の運用オブジェクトに到達できます。 これらのガードを分離境界として使用しないでください。 これらのステートメントは、テスト コードから除外し、選択した出力によって実行されるすべてのパイプライン コードから除外します。

運用上の制限事項

  • 同時実行はサポートされていません。テストとパイプラインの更新を同時に実行することはサポートされておらず、システムではそれを防ぐことはありません。 2 つの間に調整がないため、それらを同時に実行するとリソースが競合し、運用環境の更新プログラムのパフォーマンスが大幅に低下したり、テストの開始に失敗したりする可能性があります。 パイプラインが更新を実行している間はテストを開始しないでください (または、テストの実行中に更新を開始します)。テストを実行する前に、進行中の更新が完了するのを待ちます。
  • 異常終了後の一時スキーマ: 各テスト実行では、パイプラインの既定のカタログに一時スキーマ ( redirecting_<id> という名前) が作成され、実行が完了すると自動的に削除されます。 実行が異常終了した場合 (たとえば、コンピューティングが実行中に失われた場合)、実行のモック テーブルと出力テーブルを保持して、一時スキーマを残すことができます。 運用環境のデータには影響しません。 ストレージを再利用するには、パイプラインの既定のカタログで名前が redirecting_ で始まる残りのスキーマを手動で削除します。
  • テスト実行はコンピューティングを消費します。テスト実行はパイプラインのコンピューティングで実行され、通常のパイプライン更新プログラムとして課金されます。 テスト実行に対して個別の計測は行われません。
  • 完全更新はサポートされていません。選択的更新のみを使用できます。 test_pipeline.run() 選択した出力 (または選択を渡さない場合は、すべての出力) を更新します。完全更新と完全更新の選択は実装されていません。

作成と忠実性の制限

  • エディターのみの実行: テストは、Web ベースの Lakeflow パイプライン エディターから実行する必要があります。
  • Pythonテストのみ: テストはPythonで記述する必要があります。 SQL パイプラインはテストできますが、テスト自体はPythonで記述する必要があります。
  • ガバナンスの忠実性: モック データは、置き換える運用テーブルで定義されている行フィルターや列マスクを継承しません。 テスト結果は、モック入力を指定したとおりに反映され、管理された運用データに対する同じクエリの動作とは異なる場合があります。

手順 1: パイプライン設定を更新する

トリガー モードで プレビュー チャネルで実行するようにパイプラインを構成します。

  1. UI でパイプラインを開き、設定>詳細設定>チャネル>プレビュー をクリックします。
  2. パイプライン モード[トリガー済み] に設定します ([連続] を使用しないでください)。

または、パイプライン設定 JSON を直接編集します。

"continuous": false,
"channel": "PREVIEW"

手順 2: テスト ファイルを作成する

Lakeflow Pipelines エディターで、 + (追加) ボタンをクリックし、[ テスト] を選択します。 これにより、パイプラインのソース コードに含まれていないテスト ファイル (およびまだ存在しない場合は、 tests フォルダー) が作成されます。 tests フォルダーを自分で作成する必要はありません。

pytest ファイルを作成するための [テスト] オプションを示すパイプライン資産メニューを追加します。

手順 3: テストを生成する

Genie Code では、テスト スキャフォールディングを生成できます。

  • テスト ファイル内で、[テストの 生成 ] ボタンをクリックします。

    [テストの生成] ボタンが表示された空のテスト ファイル。

  • または、Genie Code エージェント モードで /tests を使用します。

    Genie Code によって TestPipeline ベースの単体テストが設定されたテスト ファイル。

Genie Code を使用して定型文を生成し、エッジ ケースに合わせてカスタマイズします。

または、テスト コードを自分で記述することもできます。 各テスト ファイルの先頭に次のインポートを追加します。

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

test_pipeline = TestPipeline.active()

手順 4: テストを実行する

Lakeflow パイプライン エディターからテストを実行します。

  • [ 再生] アイコン をクリックします。テスト関数の横にある余白にある [再生] ボタンをクリックして、個々のテストを実行します。
  • テスト ファイルの上部にある [ファイルでテストを実行 ] をクリックして、そのファイル内のすべてのテストを実行します。

テスト結果 (成功または失敗) がエディターの下部パネルに表示されます。 アサーション エラーを確認して、失敗の原因を特定します。

API のテスト

API Description
TestPipeline.active() Lakeflow パイプライン エディターで現在編集中のパイプラインの TestPipeline オブジェクトを返します。 このオブジェクトは、ソース コード、構成、既定のカタログ/スキーマなど、パイプラインへの参照です。
test_pipeline.run(test_spark, set([table_names])) パイプラインの更新を同期的に実行し、テーブル名が指定されている場合は選択的な更新を実行します。 パイプラインの実行が成功した後、または例外で終了した後に返されます。
test_spark フィクスチャ by name でテーブルを参照するテーブルの読み取りおよび書き込み(たとえば 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 メールを処理し、予期されるスキーマを生成します。

パイプライン変換:

これらの変換により、単純な 2 テーブル パイプラインが作成されます。 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 Type 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 の複数のレコードを含むモック変更フィードを作成し (更新をシミュレート)、最新のレコードのみがターゲットに保持されていることを確認します。 2 番目のテストでは、パイプラインを実行し、変更フィードにさらにイベントを追加し、パイプラインをもう一度実行することで、到着遅延イベントと順序外イベントをシミュレートします。

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 によって正しく処理されることを検証します。

パイプライン変換:

この変換により、スナップショットから自動 CDC が設定されます。スナップショット テーブルから読み取り、時間の経過に伴う変更を SCD タイプ 2 として追跡します (完全な履歴が保持されます)。

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}