在新的文件到达时触发作业

当新文件到达外部位置(例如Azure存储)时,可以使用文件到达触发器来触发作业的运行。 当计划作业的效率因不规则的新数据到达而受到影响时,此功能非常有用。

文件到达触发器的工作原理

文件到达触发器会尽力确保每分钟检查一次新文件,尽管这可能会受到底层云存储性能的影响。 除了因列出存储位置中的文件而产生的云提供商成本之外,文件到达触发器不会产生其他成本。

文件到达触发器可以配置为监视 Unity Catalog 外部位置或卷的根,或者外部位置或卷的子路径。 例如,对于 Unity 目录卷 /Volumes/mycatalog/myschema/myvolume/,下面是文件到达触发器的有效路径:

/Volumes/mycatalog/myschema/myvolume/
/Volumes/mycatalog/myschema/myvolume/mydirectory/

文件到达触发器以递归方式检查配置位置的所有子目录中是否有新文件。 例如,为位置 /Volumes/mycatalog/myschema/myvolume/mydirectory/ 创建文件到达触发器,并且此位置具有以下子目录:

/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirA
/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirB
/Volumes/mycatalog/myschema/myvolume/mydirectory/subdirC/subdirD

触发器会检查 mydirectorysubdirAsubdirBsubdirCsubdirC/subdirD 中是否有新文件。

包含文件事件的文件到达触发器

为了获得最佳性能,应为 文件事件启用外部位置。 为外部位置启用文件事件时,Azure Databricks 使用内部服务通过处理来自云提供商的更改通知来跟踪引入元数据。 此服务保留在服务确定的滚动保留期内创建或更新的最新文件的元数据,从而提高文件处理效率。

在启用外部位置的文件事件后的几分钟内,监视该外部位置涵盖路径的现有文件抵达触发器将从文件事件的启用中受益,而新的触发器将在几秒钟内受益。

有关文件事件在外部位置的性能和容量优势的详细信息,请参阅“限制”。

在您开始之前

若要使用文件到达触发器,需要满足以下要求:

添加文件到达触发器

若要将文件到达触发器添加到作业,请执行以下操作:

  1. 在 Azure Databricks 工作区的边栏中,单击Jobs & Pipelines
  2. (可选)选择作业归我所有筛选器。
  3. 单击作业的名称链接。
  4. 在右侧的“ 作业详细信息 ”窗格中,单击“ 添加触发器”。
  5. 在“触发器类型”中,选择“文件到达”
  6. 存储位置中,输入 Unity Catalog 外部位置的根或子路径的 URL,或者要监视的 Unity Catalog 卷的根或子路径的 URL。
  7. (可选)配置高级选项(两次触发之间的最短时间(秒)上次更改后的等待时间(秒)),以控制运行的触发频率。 有关设置示例,请参阅 控制触发运行的频率
  8. 若要验证配置,请单击“测试连接”。
  9. 单击“保存”。

若要稍后编辑、暂停或删除此触发器,请使用“作业详细信息”窗格的“计划和触发器”部分。 请参阅 “管理现有触发器”。

控制触发运行的频率

文件到达触发器中的两个高级选项控制文件到达事件如何触发作业运行。 这些选项采用两种常见的速率控制模式:冷却去抖

  • 触发器之间的最短时间(以秒为单位):将作业限制为按此间隔最多运行一次(运行之间的冷却时间)。 运行完成后,在冷却期内到达的文件要等到该间隔时间结束后,才会启动新一轮运行。 使用此选项可以限制创建运行的频率,以便频繁到达不会创建回退运行。
  • 最后一次更改后等待时间(秒):在最近一个文件到达后,等待这么长时间再开始运行,并且每次有新文件到达都会重置计时器(防抖)。 当文件成批到达,并且您希望在所有文件都已到达后通过一次运行处理整批文件时,请使用此选项。

可以自行设置任一选项,也可以同时设置这两个选项。 请参阅以下示例。

最多每 15 分钟运行一次

若要在文件到达时创建运行,但不超过每 15 分钟一次,请设置以下高级选项:

  • 触发器之间的最短时间(以秒为单位): 900

每次运行完成后,触发器将等待 900 秒(15 分钟),然后再启动另一个运行,即使文件继续到达。 这会将运行创建限制为每 15 分钟最多一次运行。

等待完整批次到达

当文件成批到达,并且你希望在单次运行中处理每一批文件时,请将 上次更改后等待时间(秒) 设置为一个小于批次之间时间间隔、但大于同一批次内文件之间时间间隔的值。 例如,如果新批处理大约每 5 分钟启动一次,请设置以下高级选项:

  • 等待最后一次更改(以秒为单位): 60

每个新文件都会重置计时器,因此只有在连续 60 秒没有新文件到达后,触发器才会启动一次运行。 此设置假定批内的文件在彼此的 60 秒内到达,因此计时器不会在批处理中间过期,并且批相隔 60 秒以上,因此连续批不会合并到单个运行中。

限制频率并等待完整批次完成

当你既想限制创建运行的频率,又想避免在批处理过程中途启动运行时,可以将这两个选项结合使用。 例如:

  • 触发器之间的最短时间(以秒为单位): 900
  • 等待最后一次更改(以秒为单位): 60

使用此配置,触发器在启动运行前等待批处理完成登陆(60 秒且没有新文件),并且不会每 15 分钟启动多个运行。

在到达时发现和处理文件

若要处理被文件到达触发器激活的文件,可以使用自动加载程序。 Auto Loader 以增量方式高效处理新文件,并提供精确一次处理保证。 例如,使用以下代码片段将文件加载到 Delta 表中。

若要使用此解决方案,请使用文件到达触发器创建作业,并添加包含以下代码的笔记本。 将每个 [REPLACE] 占位符替换为相应的值。

# Configuration
file_location = "[REPLACE]" # The same URL configured for the file arrival trigger.
checkpoint_location = "[REPLACE]" # a separate URL (outside `file_location`) used to store the Auto Loader checkpoint, which enables exactly-once processing.
sink_table = "[REPLACE]" # Delta table to write to

# Use Auto Loader to discover new files.
# Do not modify code below this line
streamingQuery = spark.readStream.format("cloudFiles") \
  .option("cloudFiles.format", "json") \
  .option("cloudFiles.schemaLocation", checkpoint_location) \
  .option("cloudFiles.useManagedFileEvents","true") \
  .load(file_location) \
  .writeStream \
  .option("checkpointLocation", checkpoint_location) \
  .trigger(availableNow = True) \
  .toTable(sink_table)

如果需要使用自定义逻辑处理新文件,并且只想发现新文件的 URL,则可以改用 foreachBatch ,如下面的代码片段所示。 请注意,foreachBatch 仅提供至少一次的处理保障。 请参阅 foreachBatch以获取有关如何使用 的更多信息。

# Configuration
file_location = "[REPLACE]" # The same URL configured for the file arrival trigger.
checkpoint_location = "[REPLACE]" # a separate URL (outside `file_location`) used to store the Auto Loader checkpoint, which enables exactly-once processing.

def process_batch(batch_df, batch_id):
  file_url = batch_df.select("path").collect()[0].path
  # [REPLACE] Your custom function for processing newly arrived files

# Use Auto Loader to discover new files.
# Do not modify code below this line
streamingQuery = spark.readStream.format("cloudFiles") \
  .option("cloudFiles.format", "binaryFile") \
  .option("cloudFiles.useManagedFileEvents","true") \
  .load(file_location) \
  .drop("content") \
  .writeStream \
  .foreachBatch(process_batch) \
  .option("checkpointLocation", checkpoint_location) \
  .trigger(availableNow = True) \
  .start()

接收文件到达触发器失败通知

若要在文件到达触发器评估失败时接收通知,请配置用于接收作业失败通知的电子邮件或系统目标。 请参阅为作业添加通知

限制

  • 仅运行新的文件触发器。 使用同名文件覆盖现有文件不会触发运行。
  • 如果为文件事件启用了存储位置,则向现有空文件添加内容将触发运行,因为此位置被视为新文件到达。
  • 文件事件监听 FlushWithClose 事件以处理文件。 某些Azure API 使用情况可能不会发出此事件,这可能会延迟文件发现。 若要处理此方案,请参阅 经典文件通知事件
  • 用于文件到达触发器的路径不得包含外部表或目录和架构的托管位置。

  • 用于文件到达触发器的路径不能包含通配符(例如 *?)。

  • 如果将存储位置配置为 Unity Catalog 中的外部位置,并且该外部位置已经为文件事件启用

    • 存储位置中的文件数没有限制。

    • 当存在太多多余的文件更新时,触发器可能会因超时而出错。

      在 Unity Catalog 的外部位置或卷的子路径上设置文件到达触发器时,该子路径外的更改(例如外部位置的根目录的更改)可能会增加触发器需要处理的元数据量。 在高更改环境中,这可能会导致触发器超出其处理时间限制并进入错误状态。

      若要防止这种情况,请创建一个 Unity 目录卷,该卷专门映射到要监视的子目录,并在该卷的根目录上设置文件到达触发器。 此方法将目标路径隔离为触发器的有效根目录,降低不相关的根目录层次变化的影响,并防止触发器进入错误状态。

    • 如果修改了现有文件,并且其元数据超出了滚动保留期,该修改将被视为新文件到达,从而触发作业运行。 您可以通过仅引入不可变文件来避免此问题,或者可以使用Auto Loader中的文件到达触发器来跟踪引入进度。

  • 如果未为文件事件启用存储位置:

    • 最多可以配置 50 个作业,并在 Azure Databricks 工作区中的此类位置上配置文件到达触发器。
    • 存储位置最多可以包含 10,000 个文件。 如果配置的存储位置是 Unity Catalog 外部位置或卷的子路径,则 10,000 个文件的限制适用于该子路径,而不适用于存储位置的根。 例如,存储位置的根的子目录中可以包含 10,000 多个文件,但配置的子目录不得超过 10,000 个文件的限制。

另请参阅 文件事件限制