Pruebas unitarias para canalizaciones

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 privilegios USE CATALOG y CREATE SCHEMA sobre 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 Owner del pipeline (IS OWNER); CAN RUN y CAN MANAGE no 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 CATALOG y CREATE SCHEMA. Un propietario del catálogo, un administrador del metastore o un usuario con el privilegio MANAGE puede 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 POLICY Estas 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") o df.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 ruta dbfs:/, o una ruta en la nube o en una ubicación externa como s3://... o abfss://...) 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) o spark.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.
  • 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 el event_log_table_name devuelto por la ejecución y consúltelo a través de test_spark. event_log_table_name puede ser None (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_success antes 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(incluido SET MANAGED LOCATION), GRANT/REVOKE, ALTER ... OWNER TO, SET/UNSET TAGSy .CREATE/DROP POLICY Algunos formularios SQL ejecutados mediante test_spark se 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 por redirecting_ 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.

  1. En la interfaz de usuario, abra su canalización y haga clic en Configuración>Configuración avanzada>Canal>Vista previa
  2. 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.

Añade un menú de recursos del pipeline que muestre la opción Prueba para crear un archivo pytest.

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 .

    Archivo de prueba vacío con el botón Generar pruebas.

  • Como alternativa, use /tests dentro del modo agente de Genie Code.

    Archivo de prueba rellenado por Genie Code con pruebas unitarias basadas en TestPipeline.

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 Icono de reproducció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}