Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Importante
Esta característica se encuentra en su versión beta.
Para obtener información general sobre las pruebas unitarias de Python en Databricks, consulte las pruebas unitarias de Python.
Las canalizaciones de Lakeflow admiten la escritura de pruebas unitarias Python en el Editor de canalizaciones de Lakeflow basado en web. Esto le permite validar Python o lógica de transformación de SQL mediante datos ficticios. Con el marco de pruebas de pipelines, puede probar casos extremos, validar API de pipelines propias (Auto CDC, tablas de streaming, expectativas, flujos de anexión) e iterar utilizando entradas simuladas para operaciones de identificadores de tablas compatibles. Revise las limitaciones de aislamiento antes de ejecutar pruebas.
- Ejecución de pruebas aisladas: el marco proporciona una sparkSession que redirige las operaciones de tabla a un esquema de prueba temporal en el catálogo predeterminado de la canalización, por lo que puede simular datos de entrada y escribir salidas de prueba sin afectar a las tablas de producción. El aislamiento se aplica a las operaciones que hacen referencia a una tabla por nombre; consulte Limitaciones.
- Ámbito de prueba flexible: ejecuta un subconjunto de un pipeline (tablas individuales, cadenas de tablas dependientes o pipelines completos) en el entorno de cálculo del pipeline utilizando la SparkSession de prueba.
- Validación de resultados: compruebe los resultados de las tablas de salida aisladas creadas en una prueba mediante aserciones de pytest estándar.
Cuándo usar pruebas unitarias
Entre los casos de uso típicos se incluyen:
- Validar la nueva lógica de transformación: pruebe que la transformación genera el esquema esperado, los recuentos de filas, las agregaciones y la lógica de negocios antes de ejecutarse en los datos de producción.
- Probar las especificaciones de Auto CDC: Validar que las definiciones de los flujos de Auto CDC procesen correctamente los eventos de cambio y gestionen inserciones, actualizaciones, eliminaciones y tipos SCD (dimensiones de cambio lento) con datos simulados.
- Probar las expectativas y las reglas de calidad de los datos: verificar que las expectativas fallen cuando deban y se cumplan cuando los datos sean válidos.
- Pruebas en tablas dependientes: pruebe cadenas de transformaciones (por ejemplo, bronce, plata y oro) para validar que los datos fluyen correctamente a través del gráfico de canalización.
Requirements
Permiso de pipeline
Owner, además de los privilegiosUSE CATALOGyCREATE SCHEMAsobre el catálogo predeterminado del pipeline. El marco necesita estos privilegios para crear el esquema de prueba temporal donde se ejecutan las pruebas.Para comprobar o establecer el permiso de canalización, abra la canalización y haga clic en Compartir. Debes ser el
Ownerdel pipeline (IS OWNER);CAN RUNyCAN MANAGEno son suficientes para ejecutar pruebas. Consulte Configuración de permisos de canalización.Para comprobar o establecer los privilegios de catálogo, abra el catálogo en el Explorador de catálogos, seleccione la pestaña Permisos y confirme que tiene
USE CATALOGyCREATE SCHEMA. Un propietario del catálogo, un administrador del metastore o un usuario con el privilegioMANAGEpuede concederlos, incluso mediante SQL:GRANT USE CATALOG, CREATE SCHEMA ON CATALOG <catalog_name> TO `<principal>`;Para obtener más información, consulte Referencia de privilegios del catálogo de Unity.
La canalización debe configurarse en modo activado (no continuo).
La canalización debe estar en el canal PREVIEW. Las pruebas unitarias están en beta y solo están disponibles en versión preliminar.
Spark Connect no se admite.
Note
El aislamiento de pruebas abarca las operaciones sobre tablas que hacen referencia a una tabla por su nombre. Las operaciones que omiten el aislamiento pueden producirse en el código de prueba y en cualquier código de canalización ejecutado por las salidas que seleccione, incluidas sus dependencias transitivas. Un archivo de prueba que parezca seguro puede seguir ejecutando un flujo de pipeline que lea o escriba mediante una ruta o un conector, lo que afecta a los datos de producción. Para evitar que las pruebas afecten a los datos o metadatos de producción, siga estas reglas:
- Haga referencia a cada tabla por su nombre (
catalog.schema.table) y simule todas las entradas por su nombre. No lea ni escriba mediante rutas (/Volumes/...,dbfs:/...,s3://...,abfss://...) ni lea mediante conectores como Kafka o Auto Loader. Estos omiten el aislamiento y actúan en sistemas de producción reales. - No ejecute instrucciones de gobernanza o propiedad, como
GRANT,REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSo .CREATE/DROP POLICYEstas operaciones se ejecutan en el entorno de producción real protegido. - No cree catálogos ni esquemas (
CREATE CATALOG,CREATE SCHEMA). Estas llegan a su metastore real de Unity Catalog. - No ejecute todo el proceso si su grafo incluye entradas basadas en rutas, conectores, escrituras imperativas u otros efectos secundarios externos. Seleccione solo las salidas cuyas dependencias utilicen operaciones admitidas de tablas de catálogo y que hayan sido reemplazadas por entradas simuladas.
Consulte Limitaciones para obtener más información.
Limitaciones
Warning
Algunas operaciones omiten el aislamiento de pruebas y pueden actuar en metadatos o datos de producción reales. Revise las siguientes limitaciones antes de ejecutar pruebas.
El aislamiento de las pruebas se basa únicamente en el nombre de la tabla.
No lea ni escriba mediante una ruta de acceso o un conector. El aislamiento solo redirige las operaciones que hacen referencia a una tabla por nombre (por ejemplo,
spark.read.table("catalog.schema.table")odf.write.saveAsTable("catalog.schema.table")). Las operaciones dirigidas mediante una ruta o a través de un conector eluden el aislamiento y actúan directamente sobre los sistemas de producción reales:-
La escritura mediante ruta (por ejemplo,
df.write.save("/Volumes/..."), una rutadbfs:/, o una ruta en la nube o en una ubicación externa comos3://...oabfss://...) escribe en el almacenamiento real de producción y puede sobrescribir los datos de producción. -
Leer por ruta (por ejemplo,
spark.read.load(path)ospark.read.format("delta").load(path)) devuelve datos de producción reales en lugar de tu simulacro. -
La lectura desde un conector se conecta a la fuente real de producción. Esto incluye Kafka (lee desde los brokers reales) y Auto Loader (
cloudFiles, que lee desde la ruta real de almacenamiento en la nube). Ninguna de ellas se redirige a tus datos simulados.
-
La escritura mediante ruta (por ejemplo,
No utilices la
event_log()función de valores de tabla desde una prueba unitaria de canalización. En modo de prueba,event_log()no se redirige al registro de eventos de tu ejecución de prueba. Puede devolver el registro de eventos de producción o uno registrado previamente, por lo que las aserciones sobre él podrían leer datos de producción. En su lugar, use elevent_log_table_namedevuelto por la ejecución y consúltelo a través detest_spark.event_log_table_namepuede serNone(por ejemplo, si no se puede resolver el nombre de la tabla del registro de eventos), así que compruébalo antes de consultar: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)No aserte
status.is_successantes de leer el registro de eventos si el objetivo es diagnosticar una actualización con errores. El registro de eventos suele ser lo que se inspecciona para comprender por qué se produjo un error en una actualización.
Gobernanza y operaciones DDL
- No se admiten las mutaciones de catálogo, esquema, permiso, propiedad, etiqueta y directiva. Esto incluye
CREATE/DROP/ALTER CATALOG,CREATE/DROP/ALTER SCHEMA(incluidoSET MANAGED LOCATION),GRANT/REVOKE,ALTER ... OWNER TO,SET/UNSET TAGSy .CREATE/DROP POLICYAlgunos formularios SQL ejecutados mediantetest_sparkse rechazan como defensa en profundidad; otros formularios o las mismas operaciones invocadas a través de las API directas pueden llegar a objetos de producción reales. No confíe en estos guardias como límite de aislamiento. Mantenga estas instrucciones fuera del código de prueba y de cualquier código de canalización ejecutado por las salidas seleccionadas.
Limitaciones operativas
- No se admite la ejecución simultánea: no se admite la ejecución de una prueba y una actualización de canalización al mismo tiempo y el sistema no lo impide. No hay ninguna coordinación entre los dos, por lo que ejecutarlos simultáneamente puede competir con los recursos, degradando gravemente el rendimiento de la actualización de producción o haciendo que la prueba no se inicie. No inicie una prueba mientras la canalización ejecuta una actualización (o inicia una actualización mientras se ejecuta una prueba); espere a que finalice cualquier actualización en curso antes de ejecutar pruebas.
-
Esquemas temporales después de la finalización anómala: cada ejecución de prueba crea un esquema temporal (denominado
redirecting_<id>) en el catálogo predeterminado de la canalización y lo quita automáticamente cuando finaliza la ejecución. Si una ejecución finaliza de forma anómala (por ejemplo, si se pierde el cálculo a mitad de la ejecución), el esquema temporal puede quedar almacenado, conservando las tablas de simulación y de resultados de la ejecución. No afecta a los datos de producción. Para recuperar espacio de almacenamiento, elimina manualmente cualquier esquema sobrante cuyo nombre comience porredirecting_en el catálogo predeterminado del pipeline. - Las ejecuciones de prueba consumen capacidad de cómputo: las ejecuciones de prueba se ejecutan en la capacidad de cómputo de la canalización y se facturan como actualizaciones normales de la canalización. No existe una medición independiente para las ejecuciones de prueba.
-
No se admite la actualización completa: solo está disponible la actualización selectiva.
test_pipeline.run()actualiza las salidas que selecciones (o todas las salidas si no se especifica ninguna selección); la actualización completa y la selección de actualización completa no están implementadas.
Limitaciones de edición y fidelidad
- Ejecución solo en el editor: Las pruebas deben ejecutarse desde el Editor web de Lakeflow Pipelines.
- Solo pruebas en Python: las pruebas deben escribirse en Python. Puede probar las canalizaciones de SQL, pero las propias pruebas deben escribirse en Python.
- Fidelidad de gobernanza: los datos ficticios no heredan filtros de fila ni máscaras de columna definidos en las tablas de producción que reemplaza. Los resultados de las pruebas reflejan exactamente las entradas ficticias que se proporcionan y pueden diferir de cómo se comporta la misma consulta en los datos de producción regulados.
Paso 1: Actualizar la configuración de canalización
Configura el pipeline para que se ejecute en el canal PREVIEW en modo activado.
- En la interfaz de usuario, abra su canalización y haga clic en Configuración>Configuración avanzada>Canal>Vista previa
- Establece el Modo del pipeline en Activado (no utilices Continuo).
Como alternativa, edite la configuración de canalización JSON directamente:
"continuous": false,
"channel": "PREVIEW"
Paso 2: Crear un archivo de prueba
En el Editor de canalizaciones de Lakeflow, haga clic en el + botón (agregar) y seleccione Probar. Esto crea un archivo de prueba (y la tests carpeta, si aún no existe) que no se incluye en el código fuente de la canalización. No es necesario crear la tests carpeta usted mismo.
Paso 3: Generación de pruebas
Genie Code puede generar andamiaje para pruebas:
Dentro del archivo de prueba, haga clic en el botón Generar pruebas .
Como alternativa, use
/testsdentro del modo agente de Genie Code.
Utiliza Genie Code para generar código estándar y, a continuación, personalízalo para tus casos extremos.
Como alternativa, puede escribir el código de prueba usted mismo. Agregue las siguientes importaciones a la parte superior de cada archivo de prueba:
import pytest
from pyspark.pipelines.testing import TestPipeline, test_spark
test_pipeline = TestPipeline.active()
Paso 4: Ejecutar pruebas
Ejecute pruebas desde el Editor de canalizaciones de Lakeflow:
- Haz clic en el botón
(reproducir) situado en el margen junto a una función de prueba para ejecutar una prueba individual.
- Haga clic en Ejecutar pruebas en el archivo en la parte superior del archivo de prueba para ejecutar todas las pruebas de ese archivo.
Los resultados de las pruebas (satisfactorios o fallidos) aparecen en el panel inferior del Editor. Revise los errores de aserción para depurar errores.
Probar las API
| API | Description |
|---|---|
TestPipeline.active() |
Devuelve un objeto TestPipeline de la canalización que se está editando actualmente en el Editor de canalizaciones Lakeflow. Este objeto es una referencia a la canalización, incluido su código fuente, configuraciones, catálogo/esquema predeterminado, etc. |
test_pipeline.run(test_spark, set([table_names])) |
Ejecuta sincrónicamente una actualización de la canalización, realizando una actualización selectiva si se especifican nombres de tabla. Vuelve tras la ejecución satisfactoria del pipeline o cuando este finalice con una excepción. |
test_spark fixture |
Crea una SparkSession de prueba con redirección de tablas del catálogo que redirige automáticamente las operaciones de lectura y escritura de tablas que hacen referencia, por nombre, a una tabla (por ejemplo, spark.read.table("catalog.schema.table") o df.write.saveAsTable("catalog.schema.table")) a un esquema temporal de prueba. El redireccionamiento solo se aplica a las operaciones de tabla basadas en el nombre; no cubre las lecturas ni las escrituras dirigidas mediante una ruta de acceso o a través de un conector, ya que estas actúan directamente sobre el sistema real. Consulte Limitaciones. |
Crear datos ficticios
Puede simular datos de entrada mediante SQL o 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")
Para generar volúmenes más grandes de datos sintéticos realistas, puede usar la biblioteca Faker . Ejecuta primero %pip install faker en tu pipeline y, a continuación, crea un DataFrame a partir de UDF respaldadas por 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")
Ejecuta el pipeline o tablas específicas
# 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)
Ejemplos
Ejemplo 1: Probar las agregaciones con recuento de filas, esquema y manejo de valores nulos
Objetivo: validar que la agregación de usuarios cuenta correctamente los usuarios por tipo, controla los correos electrónicos NULOs y genera el esquema esperado.
Transformaciones de canalización:
Estas transformaciones crean una canalización sencilla de dos tablas: users selecciona datos de usuario y counts agrupa usuarios por tipo y cuenta el número total de usuarios y correos electrónicos válidos.
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")
)
)
Pruebas:
Estas pruebas validan recuentos de filas, estructura de esquema, control nulo y lógica de agregación mediante la creación de datos de usuario ficticios con valores NULL intencionales y la ejecución de la canalización de forma aislada.
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)
Ejemplo 2: Prueba de CDC automático
Objetivo: Validar que Auto CDC procese correctamente la fuente de cambios con inserciones y actualizaciones.
Transformación de la canalización:
Esta transformación configura Auto CDC a partir de un flujo de cambios, que lee los cambios en streaming y los aplica a la tabla de destino como SCD de tipo 1 (mantiene solo la versión más reciente).
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
)
Pruebas:
La primera prueba crea una fuente de cambios simulada con varios registros para el mismo userId (simulando una actualización) y comprueba que solo se conserva el registro más reciente en el destino. La segunda prueba simula eventos que llegan tarde y fuera de orden ejecutando el pipeline, añadiendo más eventos al feed de cambios y volviendo a ejecutar el pipeline.
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"
Ejemplo 3: Prueba de Auto CDC a partir de una instantánea
Objetivo: valide que CDC procesa correctamente los cambios de instantánea, incluidas las inserciones, las actualizaciones y las eliminaciones.
Transformación de la canalización:
Esta transformación configura Auto CDC a partir de una instantánea, que lee desde una tabla de instantáneas y realiza un seguimiento de los cambios a lo largo del tiempo como SCD de tipo 2 (mantiene el historial completo).
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
)
Prueba:
Esta prueba crea una instantánea inicial, ejecuta la canalización y, a continuación, simula una actualización de la instantánea mediante el truncado y la inserción de nuevos datos para comprobar que CDC captura todos los cambios.
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}
Ejemplo 4: Prueba de uniones y expectativas
Objetivo: validar que las uniones funcionen correctamente y que las expectativas filtren los datos no válidos.
Transformación de canalización:
Esta transformación combina imágenes de propiedades con comodidades y aplica una expectativa de filtrar las imágenes cargadas antes de enero de 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"
)
)
Pruebas:
Estas pruebas comprueban que la combinación genera el número correcto de filas y que la expectativa filtra correctamente los registros con fechas de carga no válidas.
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}