使用控件表驱动 For each 作业

当你在多个输入上运行相同的处理,比如市场、源表、客户或日期分区时,在工作中硬编码该列表意味着每次列表变更时都要编辑代码并重新部署。 相反,将列表存储在一个 控制表 中,作业在运行时读取它。 添加或删除工作时,你更新表中的一行,下一次作业运行时会直接处理该更改,而不修改作业本身。 这是一种 元数据驱动 的模式:数据而非代码控制作业处理的内容。

本教程生成一个作业,该作业在预安装的 Wanderbricks 示例数据集上使用此模式,因此可以在不创建任何源数据的情况下端到端运行它。 情景是一个度假租赁平台,对每个 房产细分 (如 Ski ResortUrban Year-Round)进行相同的价格分析。 控制表列出要分析的段,SQL任务读取该表, For each 任务并行执行每个段一次分析。

工作原理

该工作将三项任务依次连接起来:

任务 类型 它的作用是什么
read_segments SQL 读取控制表,并将各行记录为 JSON 数组
process_segments 为每个 遍历行数组,每行启动一次嵌套任务
run_segment_analysis 笔记本或 SQL(嵌套在 For each 内) 每行运行一次,利用该行的值分析一个属性段

流程为 read_segmentsprocess_segmentsrun_segment_analysis(每行一次)。 SQL 任务的输出(即由行对象组成的 JSON 数组)通过动态值引用 {{tasks.read_segments.output.rows}} 流入 For each 任务的 Inputs 字段。 For each任务随后将每一行的字段作为参数传递给嵌套任务,参数分别为 {{input.property_type}}{{input.min_price}}和 。

先决条件

  • 有权创建作业和笔记本的Azure Databricks工作区。
  • 在 Unity Catalog 中创建表的权限,以及在目录中创建架构(即 USE CATALOGCREATE SCHEMA 权限)以容纳控制表的权限。
  • 一个SQL仓库来运行SQL任务。 如果你没有 SQL 仓库,请参见 创建 SQL 仓库
  • samples该目录在每个已启用 Unity Catalog 的工作区中均可用。 本教程从 samples.wanderbricks.properties 读取,因此无需设置源数据。

步骤 1:创建控件表

控制表是你的作业所处理的分段列表的权威来源。 要更改作业的行为,应更新这个表,而不是作业本身。

在Azure Databricks笔记本或 SQL 编辑器中运行以下 SQL。 第一个语句创建一个用于存储控制表的模式,第二个语句创建了每个物业段一行的表,并列出该段分析中包含的最低挂牌价:

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);

用一个可以创建模式的目录来替代 <catalog-name> ,比如你的工作区目录。 在教程引用 config.property_segments的所有地方使用相同的目录,包括第三步的查找查询。

完成此步骤后,config.property_segments 包含三行,每个段对应一行。 每一行都包含作业传递给每次迭代的两个值:要分析的 property_type 和用于筛选的下限值 min_price

步骤 2:编写分析逻辑

For each 任务中的嵌套任务会针对控制表的每一行运行一次,并接收该行的 property_typemin_price 作为参数。 你可以把这个逻辑写成笔记本任务,也可以写成SQL任务。 根据你的商业逻辑做出选择:

  • 当每次迭代逻辑需要过程代码、多语言或库(例如数据科学或机器学习步骤)时,使用 笔记本任务
  • 当逻辑是单一查询或转换,且你可以声明式表达时,使用 SQL任务 。 SQL 任务需要 SQL 仓库。

以下两种变体得出相同的结果:针对被处理的细分市场,显示其价格下限或以上的挂牌数量及其平均价格。

笔记本任务

在路径中创建新笔记本,例如/Workspace/Users/<username>/run_segment_analysis。 该笔记本在 For each 任务的每次迭代中运行一次,每次接收不同的片段。

将以下代码添加到笔记本:

# 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)

注释

dbutils.widgets.text()之前调用dbutils.widgets.get()。 如果先调用 get,在作业外运行笔记本会引发 InputWidgetNotDefined 错误。

SQL 任务

SQL任务运行保存查询,所以现在就在SQL编辑器中创建并保存分析查询。 在第4步配置 For each 任务时,你会把它附加到嵌套任务上。

  1. 在你的 Azure Databricks 工作区中,点击加号图标。>查询图标。查询以打开SQL编辑器。

  2. 输入以下查询。 SQL 任务使用 :param_name 语法引用参数,因此查询会从 :property_type: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. 点击你SQL文件标签页标题中的标题 New Query <date> ,并给它命名 run_segment_analysis。 然后点击 “保存 ”,把它移到你想存放的文件夹里。

For each 任务在运行时将每次迭代的值传递给名为 :property_type:min_price 的参数。 与笔记本小部件不同,SQL命名参数不支持默认值:如果未传递参数,查询将失败,且参数解析错误。

步骤 3:创建查找查询

查找任务通过保存的查询读取控制表。 和步骤2一样,现在在SQL编辑器中创建并保存查询,然后在步骤4中附加到查找任务中。

  1. 在你的 Azure Databricks 工作区中,点击加号图标。>查询图标。查询以打开SQL编辑器。

  2. 请输入以下内容,使用您在第一步中选择的同一目录:

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

    名称是完全限定的,因为运行该查询的 SQL 仓库可能默认使用与你创建表的目录不同的目录。

  3. 点击你SQL文件标签页标题中的标题 New Query <date> ,并给它命名 read_segments。 然后点击 “保存 ”,把它移到你想存放的文件夹里。

步骤 4:创建和配置作业

保存好两个查询后,创建作业并添加两个任务:读取控制表的SQL查找任务和 For each 执行每行分析的任务。

创建任务

在 Azure Databricks 工作区中,在侧边栏中单击 加号图标。新建>工作流图标。作业。 给作业起一个描述性名称,比如 Segment Analysis

配置 SQL 查找任务

该任务读取控制表,并通过运行你在步骤 3 中保存的 read_segments 查询,使其行可供 For each 任务使用。

  1. 点击 SQL 查询 图块来配置第一个任务。 如果 SQL 查询图块不可用,请点击 添加其他任务类型,然后搜索 SQL 查询
  2. 任务名称 设置为 read_segments.
  3. 如有必要,请从类型下拉菜单中选择SQL查询
  4. SQL查询 字段中,选择 read_segments 你在步骤3中保存的查询。
  5. SQL 仓库 设置为工作区中的仓库。
  6. 单击“创建任务”。

当该任务运行时,Azure Databricks 会将结果作为 JSON 数组在 tasks.read_segments.output.rows中捕获。 SQL 任务输出总是以 JSON 数组形式返回,所以不需要额外的配置。 引用的一般形式是 tasks.<task-name>.output.rows,其中 <task-name> 与你设置的任务名称相匹配。 输出如下所示:

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

配置 For each 任务

For each 任务读取 SQL 输出并启动每个行的一个嵌套任务运行。

  1. 点击加号图标。添加任务并选择“针对每个”。

  2. 任务名称 设置为 process_segments.

  3. 确认 Depends on 已设置为 read_segments

  4. 输入 字段中,输入SQL任务捕获的行数组:

    {{tasks.read_segments.output.rows}}
    
  5. 并发2设置为并行运行两次迭代。 当嵌套任务支持更高的并行度时增加此值。

  6. 要完成此任务,请点击 添加要循环执行的任务,然后配置在每次迭代时运行的嵌套任务。

For each 任务及其嵌套任务会作为单个任务一并创建。 根据你在步骤2中选择的类型配置嵌套任务:

笔记本任务

  1. 任务名称 设置为 run_segment_analysis.

  2. “类型 ”设置为 “笔记本”。

  3. 设置 路径 到你在步骤2创建的笔记本。

  4. 点击 参数,然后点击 添加 以添加每个参数:

    • property_type{{input.property_type}}
    • min_price{{input.min_price}}

    每个 {{input.<key>}} 引用解析为当前迭代行中的匹配字段。

  5. 点击 创建任务 ,将该任务及其嵌套任务一起创建 For each

SQL 任务

该任务执行你在第 2 步中保存的 run_segment_analysis 查询。

  1. 任务名称 设置为 run_segment_analysis.

  2. 类型 设为 SQL,然后把 SQL 任务 设为 Query

  3. SQL查询 字段中,选择 run_segment_analysis 你在步骤2中保存的查询。

  4. SQL 仓库 设置为工作区中的仓库。

  5. 点击 参数,然后点击 添加 以添加每个参数:

    • property_type{{input.property_type}}
    • min_price{{input.min_price}}

    每个 {{input.<key>}} 引用解析为当前迭代行中的匹配字段。

  6. 点击 创建任务 ,将该任务及其嵌套任务一起创建 For each

你的作业的有向无环图(DAG)现在显示为 read_segments 流向 process_segments,嵌套任务位于 For each 节点内。

步骤 5:运行作业并验证

  1. 单击“ 立即运行 ”以触发作业。
  2. 选择 “跑步” 标签查看跑步。 作业的第一次运行需要几分钟才能开始计算;完成后,它会出现在列表中。
  3. 点击 process_segments 节点以展开 For each 任务。
  4. 运行页面显示一个迭代表格,每个分段对应一行,每行显示其状态、开始时间和持续时间。
  5. 点击任意迭代行即可打开其输出并确认已分析预期段。

你可以独立查看每次迭代的结果。 如果某个迭代失败,你可以只从作业运行页面重跑该迭代,而不重跑整个作业。

扩展图案

要将一段添加到分析中,请在控制表中插入一行:

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

下一次作业运行包含新段,且不更改作业配置或笔记本编辑。

同样的模式适用于任何需要数据驱动迭代的情况:

  • 按客户处理:每个客户 ID 对应一行。 嵌套任务会执行针对特定客户的转换,或将内容传送到特定客户的目标位置。
  • 表引入:每个源表名称一行。 嵌套任务读取并记录每个表格。
  • 回填处理:每个日期分区对应一行。 嵌套任务会重新处理该分区的历史数据。
  • 基于功能标志的执行:每个已启用的功能或实验各占一行。 嵌套任务激活相应的逻辑。

要停止处理某行而不删除它,可以在控制表中添加你自己的列(比如 active 标志),并在SQL查找任务中对它进行过滤。 这是一个普通的列,你可以定义并填充;这个 For each 任务本身没有内置的概念。 先添加列,然后将现有行设置为:TRUE

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

然后在 read_segments 查询中按其进行筛选,以便只有处于活动状态的行参与迭代:

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

其他资源