Uso de una tabla de control para controlar un For each trabajo

Cuando ejecutas el mismo procesamiento en muchas entradas, como mercados, tablas fuente, clientes o particiones de fecha, codificar esa lista en tu trabajo implica editar código y volver a desplegarla cada vez que cambia la lista. En su lugar, almacena la lista en una tabla de control que el trabajo lee en tiempo de ejecución. Para añadir o eliminar trabajo, actualizas una fila en la tabla y la siguiente ejecución del trabajo recoge el cambio sin editar el trabajo en sí. Este es un patrón impulsado por metadatos : los datos, no el código, controlan lo que procesa el trabajo.

Este tutorial construye un trabajo que utiliza este patrón en el conjunto de datos de ejemplo preinstalado de Wanderbricks, para que puedas ejecutarlo de principio a fin sin crear ningún dato fuente. El escenario es una plataforma de alquiler vacacional que realiza el mismo análisis de precios para cada segmento de propiedad (como Ski Resort o Urban Year-Round). Una tabla de control lista los segmentos a analizar, una tarea SQL lee esa tabla y una For each tarea ejecuta el análisis una vez por segmento, en paralelo.

Cómo funciona

El trabajo conecta tres tareas en secuencia:

tarea Tipo Qué hace
read_segments SQL Lee la tabla de control y captura las filas como un array JSON
process_segments For Each Itera sobre el array de filas, lanzando la tarea anidada una vez por fila
run_segment_analysis Cuaderno o SQL (anidado en su interior For each) Se ejecuta una vez por fila, usando los valores de esa fila para analizar un segmento de propiedad

El flujo es read_segmentsprocess_segmentsrun_segment_analysis (una vez por fila). La salida de la tarea SQL, un array JSON de objetos fila, fluye hacia el For each campo Entradas de la tarea a través de la referencia {{tasks.read_segments.output.rows}}de valor dinámico . La For each tarea pasa entonces los campos de cada fila a la tarea anidada como parámetros, disponibles como {{input.property_type}} y {{input.min_price}}.

Prerrequisitos

  • Un espacio de trabajo de Azure Databricks con permiso para crear trabajos y cuadernos.
  • Permiso para crear tablas en el Catálogo de Unity, y permiso para crear un esquema en un catálogo (los USE CATALOG privilegios y CREATE SCHEMA de Y) para contener la tabla de control.
  • Un almacén SQL para ejecutar las tareas SQL. Si no tienes uno, consulta Crear un almacén SQL.
  • El samples catálogo, que está disponible en todos los espacios de trabajo habilitados para el Catálogo de Unity. El tutorial lee de samples.wanderbricks.properties, así que no hay datos fuente que configurar.

Paso 1: Crear la tabla de control

La tabla de control es la fuente de verdad para la lista de segmentos que procesa tu trabajo. Para cambiar lo que hace el trabajo, actualizas esta tabla, no el trabajo.

Ejecuta el siguiente SQL en un cuaderno de Azure Databricks o en el editor SQL. La primera instrucción crea un esquema para contener la tabla de control, y la segunda crea la tabla con una fila por segmento de propiedad y el precio mínimo de venta que debe incluir en el análisis de ese segmento:

USE CATALOG <catalog-name>;

CREATE SCHEMA IF NOT EXISTS config;

CREATE OR REPLACE TABLE config.property_segments AS
SELECT * FROM VALUES
  ('Urban Year-Round', 150),
  ('Summer Getaway', 200),
  ('Ski Resort', 250)
AS t(property_type, min_price);

Sustituye <catalog-name> por un catálogo en el que puedas crear esquemas, como el catálogo de tu espacio de trabajo. Usa el mismo catálogo en todos los sitios que hace referencia config.property_segmentsel tutorial, incluyendo la consulta de búsqueda en el Paso 3.

Después de este paso, config.property_segments contiene tres filas, una por segmento. Cada fila lleva los dos valores que el trabajo pasa a cada iteración: el property_type para analizar y el min_price suelo para filtrar.

Paso 2: Escribe la lógica de análisis

La tarea anidada dentro de la For each tarea se ejecuta una vez por fila de la tabla de control, recibiendo los parámetros de property_type esa fila y min_price como parámetros. Puedes escribir esta lógica como una tarea de notebook o una tarea SQL. Elige según la lógica de tu negocio:

  • Utiliza una tarea de notebook cuando la lógica de iteración por iteración necesite código procedimental, múltiples lenguajes o librerías (por ejemplo, un paso de ciencia de datos o aprendizaje automático).
  • Usa una tarea SQL cuando la lógica sea una única consulta o transformación que puedas expresar de forma declarativa. Una tarea SQL necesita un almacén SQL.

Ambas variantes a continuación producen el mismo resultado: para el segmento en proceso, el número de anuncios a su precio mínimo o superior y su precio medio.

Tarea de cuaderno

Cree un nuevo cuaderno en una ruta como /Workspace/Users/<username>/run_segment_analysis. Este cuaderno se ejecuta una vez por cada iteración de la For each tarea, recibiendo un segmento diferente cada vez.

Agregue el código siguiente al cuaderno:

# Set default values so you can run the notebook on its own while developing.
# When the notebook runs inside a For each task, the job overrides these defaults.
dbutils.widgets.text("property_type", "Ski Resort", "Property type")
dbutils.widgets.text("min_price", "250", "Minimum price")

# Read the parameters passed by the For each task.
property_type = dbutils.widgets.get("property_type")
min_price = dbutils.widgets.get("min_price")

result = spark.sql(
    """
    SELECT :property_type AS property_type,
           COUNT(*) AS property_count,
           ROUND(AVG(base_price), 2) AS avg_price
    FROM samples.wanderbricks.properties
    WHERE property_type = :property_type
      AND base_price >= :min_price
    """,
    args={"property_type": property_type, "min_price": min_price},
)
display(result)

Note

Llame a dbutils.widgets.text() antes de dbutils.widgets.get(). Si llamas get primero, ejecutar el portátil fuera de un trabajo genera un InputWidgetNotDefined error.

Tarea SQL

Una tarea SQL ejecuta una consulta guardada, así que crea y guarda la consulta de análisis en el editor SQL ahora. Lo asocias a la tarea anidada cuando configuras la For each tarea en el Paso 4.

  1. En tu espacio de trabajo de Azure Databricks, haz clic en icono Plus.Nuevo>Icono de consulta.Consulta para abrir el editor SQL.

  2. Escriba la siguiente consulta. Las tareas SQL hacen referencia a parámetros con la :param_name sintaxis, por lo que la consulta lee su segmento y el piso de precio a partir de los :property_type parámetros y::min_price

    SELECT :property_type AS property_type,
           COUNT(*) AS property_count,
           ROUND(AVG(base_price), 2) AS avg_price
    FROM samples.wanderbricks.properties
    WHERE property_type = :property_type
      AND base_price >= :min_price;
    
  3. Haz clic en el título New Query <date> en el encabezado de la pestaña de tu archivo SQL y ponle el nombre run_segment_analysis. Luego haz clic en Guardar para moverlo a una carpeta donde quieres guardarlo.

La For each tarea pasa los valores de cada iteración a los :property_type parámetros y :min_price nombrados en tiempo de ejecución. A diferencia de los widgets de cuaderno, los parámetros con nombre SQL no soportan valores por defecto: si un parámetro no se pasa, la consulta falla con un error de resolución de parámetros.

Paso 3: Crear la consulta de búsqueda

La tarea de búsqueda lee la tabla de control a través de una consulta guardada. Como en el Paso 2, crea y guarda la consulta en el editor SQL ahora, y luego addícala a la tarea de búsqueda en el Paso 4.

  1. En tu espacio de trabajo de Azure Databricks, haz clic en icono Plus.Nuevo>Icono de consulta.Consulta para abrir el editor SQL.

  2. Introduce lo siguiente, usando el mismo catálogo que elegiste en el Paso 1:

    SELECT property_type, min_price FROM <catalog-name>.config.property_segments;
    

    El nombre está totalmente condicionado porque el almacén SQL que ejecuta esta consulta puede recurrir por defecto a un catálogo diferente al que creaste la tabla.

  3. Haz clic en el título New Query <date> en el encabezado de la pestaña de tu archivo SQL y ponle el nombre read_segments. Luego haz clic en Guardar para moverlo a una carpeta donde quieres guardarlo.

Paso 4: Crear y configurar el trabajo

Con ambas consultas guardadas, crea el trabajo y añade sus dos tareas: la tarea de consulta SQL que lee la tabla de control y la For each tarea que ejecuta el análisis de cada fila.

Crear el puesto

En tu espacio de trabajo de Azure Databricks, en la barra lateral haz clic en icono Plus.Nuevo>Icono de flujos de trabajo.Trabajo. Ponle al puesto un nombre descriptivo, como Segment Analysis.

Configurar la tarea de consulta SQL

Esta tarea lee la tabla de control y pone sus filas a disposición For each ejecutando la read_segments consulta que guardaste en el Paso 3.

  1. Haz clic en la tesela de consulta SQL para configurar la primera tarea. Si la tesela de consulta SQL no está disponible, haz clic en Añadir otro tipo de tarea y busca consulta SQL.
  2. Establezca Nombre de tarea en read_segments.
  3. Si es necesario, selecciona consulta SQL en el menú desplegable de Tipo .
  4. En el campo de consulta SQL , selecciona la read_segments consulta que guardaste en el Paso 3.
  5. Establezca SQL Warehouse para un almacén en su área de trabajo.
  6. Haga clic en Create task (Crear tarea).

Cuando esta tarea se ejecuta, Azure Databricks captura el resultado como un array JSON en tasks.read_segments.output.rows. La salida de la tarea SQL siempre se devuelve como un array JSON, así que no necesitas ninguna configuración adicional. La forma general de la referencia es tasks.<task-name>.output.rows, donde <task-name> coincide con el nombre de la tarea que has asignado. La salida es similar a esta:

[
  { "property_type": "Urban Year-Round", "min_price": 150 },
  { "property_type": "Summer Getaway", "min_price": 200 },
  { "property_type": "Ski Resort", "min_price": 250 }
]

Configurar la For each tarea

La tarea For each lee la salida de SQL e inicia una ejecución de tarea anidada por tabla de origen.

  1. Haz clic en el icono de más. Añadir tarea y seleccionar Para cada una.

  2. Establezca Nombre de tarea en process_segments.

  3. Verifica que Depends on está configurado como read_segments.

  4. En el campo Entradas , introduce el array de filas capturado por la tarea SQL:

    {{tasks.read_segments.output.rows}}
    
  5. Configura la concurrencia para 2 ejecutar dos iteraciones en paralelo. Aumente este valor cuando la tarea anidada admita un paralelismo superior.

  6. Para completar esta tarea, haz clic en Añadir una tarea para que se repita y configura la tarea anidada que se ejecuta en cada iteración.

La For each tarea y su tarea anidada se crean juntas como una sola tarea. Configura la tarea anidada según el tipo que elegiste en el Paso 2:

Tarea de cuaderno

  1. Establezca Nombre de tarea en run_segment_analysis.

  2. Establezca Tipo en Notebook.

  3. Establece la ruta en el cuaderno que creaste en el Paso 2.

  4. Haz clic en Parámetros y luego haz clic en Añadir para añadir cada parámetro:

    • Clave: property_type, Valor: {{input.property_type}}
    • Clave: min_price, Valor: {{input.min_price}}

    Cada {{input.<key>}} referencia se resuelve al campo correspondiente de la fila de la iteración actual.

  5. Haz clic en Crear tarea para crear la For each tarea y su tarea anidada juntas.

Tarea SQL

Esta tarea ejecuta la run_segment_analysis consulta que guardaste en el Paso 2.

  1. Establezca Nombre de tarea en run_segment_analysis.

  2. Ponle Tipo en SQL y luego pon la tarea SQL en Consulta.

  3. En el campo de consulta SQL , selecciona la run_segment_analysis consulta que guardaste en el Paso 2.

  4. Establezca SQL Warehouse para un almacén en su área de trabajo.

  5. Haz clic en Parámetros y luego haz clic en Añadir para añadir cada parámetro:

    • Clave: property_type, Valor: {{input.property_type}}
    • Clave: min_price, Valor: {{input.min_price}}

    Cada {{input.<key>}} referencia se resuelve al campo correspondiente de la fila de la iteración actual.

  6. Haz clic en Crear tarea para crear la For each tarea y su tarea anidada juntas.

Tu Grafo Acíclico Dirigido (DAG) de tarea ahora muestra read_segments fluyendo en process_segments, con la tarea anidada dentro del For each nodo.

Paso 5: Ejecuta el trabajo y verifica

  1. Haga clic en Ejecutar ahora para desencadenar el trabajo.
  2. Selecciona la pestaña Carreras para ver la carrera. La primera ejecución de un trabajo tarda unos minutos en comenzar el cálculo; Cuando termina, aparece en la lista.
  3. Haz clic en el process_segments nodo para ampliar la For each tarea.
  4. La página de la carrera muestra una tabla de iteraciones, una fila por segmento, cada una con su estado, hora de inicio y duración.
  5. Haz clic en cualquier fila de iteración para abrir su salida y confirmar que ha analizado el segmento esperado.

Puedes ver los resultados de cada iteración de forma independiente. Si falla una iteración específica, solo puedes reejecutar esa iteración desde la página de ejecución del trabajo sin volver a ejecutar todo el trabajo.

Extender el patrón

Para añadir un segmento al análisis, inserta una fila en la tabla de control:

INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);

La siguiente ejecución de trabajo incluye el nuevo segmento, sin cambios en la configuración del trabajo ni ediciones en el cuaderno.

Este mismo patrón funciona en cualquier caso en el que quieras que los datos impulsen la iteración:

  • Procesamiento por cliente: una fila por identificador de cliente. La tarea anidada aplica transformaciones específicas del cliente o entrega a destinos específicos del cliente.
  • Ingesta de tablas: una fila por nombre de tabla fuente. La tarea anidada lee e ingiere cada tabla.
  • Procesamiento de reposición: una fila por partición de fecha. La tarea anidada reprocesa datos históricos de esa partición.
  • Ejecución basada en banderas de funcionalidad: una fila por cada funcionalidad o experimento habilitados. La tarea anidada activa la lógica correspondiente.

Para dejar de procesar una fila sin eliminarla, añade tu propia columna a la tabla de control (como una active bandera) y filtra en la tarea de consulta SQL. Esta es una columna ordinaria que defines y llenas; La For each tarea no tiene un concepto incorporado. Primero añade la columna, luego establece las filas existentes en TRUE:

ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;

Luego filtra en la read_segments consulta para que solo las filas activas impulsen la iteración:

SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;

Recursos adicionales