Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Important
Эта функция доступна в бета-версии.
Общие сведения о модульном тестировании Python в Databricks см. в разделе Модульное тестирование Python.
Конвейеры Lakeflow поддерживают создание модульных тестов на Python в веб-редакторе Lakeflow Pipelines. Это позволяет проверить логику преобразования Python или SQL с помощью макетных данных. С помощью платформы тестирования конвейеров можно протестировать пограничные варианты, проверить собственные API конвейера (Auto CDC, потоковую передачу таблиц, ожидания, добавочные потоки) и выполнить итерацию с помощью макетных входных данных для поддерживаемых операций с идентификатором таблицы. Просмотрите ограничения изоляции перед выполнением тестов.
- Изолированное выполнение теста. Платформа предоставляет SparkSession, которая перенаправляет операции таблицы во временную схему тестирования в каталоге по умолчанию конвейера, поэтому вы можете макетировать входные данные и записывать выходные данные теста без влияния на рабочие таблицы. Изоляция применяется к операциям, ссылающимся на таблицу по имени; см. ограничения.
- Гибкая область теста: выполнение подмножества конвейера (отдельные таблицы, цепочки зависимых таблиц или целых конвейеров) на вычислительных ресурсах конвейера с помощью тестовой службы SparkSession.
- Проверка результатов. Проверьте результаты изолированных выходных таблиц, созданных в тесте с помощью стандартных утверждений pytest.
Когда следует использовать модульное тестирование
К типичным сценариям использования относятся:
- Проверка новой логики преобразования: Убедитесь, что преобразование формирует ожидаемую схему, количество строк, агрегаты и бизнес-логику перед запуском на рабочих данных.
- Тестирование спецификаций Auto CDC: убедитесь, что определения потока Auto CDC правильно обрабатывают события изменения, обрабатывают вставки, обновления, удаления и SCD (медленно меняющееся измерение) с помощью макетных данных.
- Тестирование ожиданий и правил качества данных: убедитесь, что ожидания не выполняются, когда это ожидается, и успешно выполняются, когда данные валидны.
- Тестирование в связанных таблицах: тестируйте цепочки преобразований (например, bronze, silver и gold), чтобы убедиться, что данные корректно проходят по графу конвейера.
Requirements
Разрешение конвейера
Owner, а такжеUSE CATALOGCREATE SCHEMAпривилегии в каталоге конвейера по умолчанию. Платформа нуждается в этих привилегиях для создания временной схемы тестирования, в которой выполняются тесты.Чтобы проверить или задать разрешение конвейера, откройте конвейер и нажмите кнопку "Общий доступ". Требуется конвейер
Owner(IS OWNER);CAN RUNиCAN MANAGEнедостаточны для запуска тестов. См. раздел "Настройка разрешений конвейера".Чтобы проверить или задать привилегии каталога, откройте каталог в обозревателе каталогов, перейдите на вкладку "Разрешения" и убедитесь, что у вас есть
USE CATALOG.CREATE SCHEMAВладелец каталога, администратор метахранилища или пользователь с привилегиейMANAGEмогут предоставить эти привилегии, в том числе с помощью SQL:GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;Дополнительные сведения см. в справочнике по привилегиям каталога Unity.
Конвейер должен быть настроен в режиме запуска (не непрерывном).
Конвейер должен находиться в канале PREVIEW . Модульное тестирование находится в бета-версии и доступно только в предварительной версии.
Spark Connect не поддерживается.
Note
Проверка изоляции охватывает операции с таблицами, ссылающиеся на таблицу по имени. Операции, которые обходят изоляцию, могут возникать как в тестовом коде, так и в любом коде конвейера, выполняемом выбранными выходными данными, включая транзитивные зависимости. Даже тестовый файл, который выглядит безопасным, может запустить процесс конвейера, который читает или записывает данные через путь или коннектор и при этом работает с производственными данными. Чтобы тесты не влияли на производственные данные или метаданные, соблюдайте следующие правила:
- Ссылка на каждую таблицу по имени (
catalog.schema.table) и макетирование всех входных данных по имени. Не выполняйте чтение или запись по пути (/Volumes/...,dbfs:/...,s3://...,abfss://...) и не выполняйте чтение из коннекторов, таких как Kafka или Auto Loader. Они обходят изоляцию и воздействуют на реальные продукционные системы. - Не выполняйте инструкции управления или владения, например
GRANT, ,REVOKE,ALTER ... OWNER TOSET/UNSET TAGSили .CREATE/DROP POLICYОни выполняются на реальном защищаемом объекте в производственной среде. - Не создавайте каталоги или схемы (
CREATE CATALOG,CREATE SCHEMA). Они обращаются к реальному хранилищу метаданных Unity Catalog. - Не запускайте весь конвейер, если в его граф входят входные данные, задаваемые путями, коннекторы, императивные записи или другие внешние побочные эффекты. Выберите только выходные данные, зависимости которых используют поддерживаемые операции с таблицами каталога и заменены макетными входными данными.
Дополнительные сведения см. в разделе Ограничения.
Ограничения
Предупреждение
Некоторые операции обходят тестовую изоляцию и могут работать с реальными производственными данными или метаданными. Ознакомьтесь со следующими ограничениями перед выполнением тестов.
Изоляция тестов — только по имени таблицы
Не читайте и не записывайте, используя путь или соединитель. Изоляция перенаправляет только операции, ссылающиеся на таблицу по имени (например,
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 (считывающий данные из реальных брокеров) и Auto Loader (
cloudFiles, считывающий данные из реального пути в облачном хранилище). Также не перенаправляется на ваши тестовые данные.
-
Запись по пути (например, по пути
Не используйте табличную функцию
event_log()в модульном тесте конвейера. В тестовом режимеevent_log()не перенаправляется в журнал событий тестового запуска. Он может возвращать рабочий или ранее зарегистрированный журнал событий, поэтому утверждения против него могут считывать рабочие данные. Вместо этого используйтеevent_log_table_nameвозвращаемый запуском и запросите его черезtest_spark.event_log_table_nameможет бытьNone(например, если имя таблицы журнала событий не может быть разрешено), поэтому проверьте его перед запросом: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 LOCATIONGRANT/REVOKE), ,ALTER ... OWNER TOSET/UNSET TAGS, и .CREATE/DROP POLICYНекоторые формы SQL, выполняемые черезtest_spark, отклоняются в рамках принципа глубокоэшелонированной защиты; другие формы, а также те же операции, вызываемые через прямые API, могут получать доступ к реальным производственным объектам. Не следует полагаться на эти защитные механизмы как на границу изоляции. Не включайте эти инструкции в тестовый код и в любой код конвейера, выполняемый при использовании выбранных выходных данных.
Операционные ограничения
- Параллельное выполнение не поддерживается: выполнение теста и обновления конвейера одновременно не поддерживается, и система не предотвращает ее. Между ними нет координации, поэтому одновременное выполнение их может бороться за ресурсы, серьезно ухудшая производительность рабочего обновления или вызывая неудачный запуск теста. Не запускайте тест во время выполнения конвейера обновления (или запускайте обновление во время выполнения теста); Дождитесь завершения любого выполняемого обновления перед выполнением тестов.
-
Временные схемы после ненормального завершения: каждый тестовый запуск создает временную схему (именованную
redirecting_<id>) в каталоге по умолчанию конвейера и автоматически удаляет ее при завершении выполнения. Если выполнение завершается аварийно (например, если во время выполнения теряется вычислительный ресурс), временная схема может остаться, содержащая имитационные и выходные таблицы этого выполнения. Это не влияет на рабочие данные. Чтобы восстановить хранилище, вручную удалите все оставшиеся схемы, имена которых начинаются сredirecting_каталога конвейера по умолчанию. - Тестовые запуски используют вычислительные ресурсы: тестовые запуски выполняются на вычислительных ресурсах конвейера и выставляются как обычные обновления конвейера. Для тестовых запусков нет отдельного измерения.
-
Полное обновление не поддерживается: доступно только выборочное обновление.
test_pipeline.run()обновляет выбранные результаты (или все результаты, если не указать выбор); полное обновление и выбор с полным обновлением не реализованы.
Ограничения разработки и точности
- Выполнение только в редакторе: тесты необходимо запускать в веб-редакторе конвейеров Lakeflow.
- Только тесты на Python: Тесты должны быть написаны на Python. Вы можете протестировать конвейеры SQL, но сами тесты должны быть записаны в Python.
- Точность управления. Макет данных не наследует фильтры строк или маски столбцов, определенные в рабочих таблицах, которые он заменяет. Результаты теста отражают макетные входные данные точно так же, как они предоставляются и могут отличаться от того, как тот же запрос работает с управляемыми рабочими данными.
Шаг 1. Обновление параметров конвейера
Настройте конвейер для запуска в канале PREVIEW в активированном режиме.
- В интерфейсе пользователя откройте свой конвейер и выберите Параметры>Дополнительные параметры>Канал>Предварительная версия
- Установите для параметра Pipeline mode значение Triggered (не используйте значение Continuous).
Кроме того, измените параметры конвейера в формате JSON напрямую:
"continuous": false,
"channel": "PREVIEW"
Шаг 2. Создание тестового файла
В редакторе конвейеров Lakeflow нажмите кнопку + (добавить) и выберите Тест. Это создает тестовый файл (и папку tests, если она еще не существует), которые не включаются в исходный код вашего конвейера. Вам не нужно создавать папку tests самостоятельно.
Шаг 3. Создание тестов
Genie Code может создавать шаблон тестирования:
В тестовом файле нажмите кнопку "Создать тесты ".
Кроме того, используйте
/testsв режиме агента 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 приспособление |
Создает тестовую версию SparkSession с перенаправлением таблицы каталога, которая автоматически перенаправляет операции чтения и записи, ссылающиеся на таблицу по имени (например, spark.read.table("catalog.schema.table") или df.write.saveAsTable("catalog.schema.table")) во временную схему теста. Перенаправление применяется только к операциям с таблицами по именам; оно не распространяется на чтение или запись по пути либо через коннектор, которые выполняются непосредственно в реальной системе. См. Ограничения. |
Создание макетных данных
Вы можете имитировать входные данные с помощью 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 в вашем конвейере, затем создайте DataFrame из UDF, поддерживаемых Faker:
# 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
Цель. Проверка агрегирования пользователей правильно подсчитывает пользователей по типу, обрабатывает пустые сообщения электронной почты и создает ожидаемую схему.
Преобразования конвейера:
Эти преобразования создают простой конвейер с двумя таблицами: 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: Тестирование Auto CDC
Цель: Проверить, что Auto CDC правильно обрабатывает поток изменений с операциями вставки и обновления.
Преобразование конвейера:
Это преобразование настраивает Auto 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 (имитация обновления) и проверяет, что только последняя запись сохраняется в целевом объекте. Второй тест имитирует поздно поступающие события и события, поступающие не по порядку, запуская конвейер, добавляя больше событий в ленту изменений и запуская его снова.
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 правильно обрабатывает изменения снимка, включая вставки, обновления и удаления.
Преобразование конвейера:
Это преобразование настраивает Auto CDC на основе снимка: Auto CDC считывает данные из таблицы снимков и отслеживает изменения с течением времени как медленно изменяющееся измерение типа 2 (SCD Type 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 года.
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}