Auto Loader 最佳实践

本页介绍了一些最佳实践,你可以用它们来配置 Auto Loader,使其能够针对你的使用场景以可靠、经济高效且可扩展的方式运行。

这些最佳实践可降低运维开销,并避免在生产环境中常见但难以诊断的问题,例如:因全量目录扫描而产生的不必要的 LIST API 成本、架构漂移导致的静默数据丢失,以及检查点配置错误导致的管道重启。

有关生产配置详细信息,请参阅 为生产工作负荷配置自动加载程序。 有关监视和可观测性,请参阅 “监视并观察自动加载程序”。

选择正确的执行框架

用例的最佳执行框架取决于对管道所需的控制程度以及要管理的运营开销。 对于大多数用户和生产流水线而言,将 Auto Loader 与 Lakeflow pipelines 结合使用是一个不错的选择。 但是,如果需要最大控制和自定义,请使用带结构化流式处理的自动加载程序。 对于具有托管体验的最简单设置,请使用托管的 LakeFlow 连接器(如果可用)。

Lakeflow 管道通过增加自动扩缩容、数据质量检查、架构演进处理以及通过事件日志进行监控等功能,扩展了结构化流。 Databricks 建议针对大多数生产引入工作负荷使用 Lakeflow 管道。

选择合适的调度和触发器类型

用例的最佳计划和触发器类型取决于延迟要求和文件到达模式。 对于大多数用例,Databricks 建议使用启用了文件事件的文件到达触发器。 这可实现低成本的低延迟引入,因为计算仅在新文件到达时运行。 这三种触发器类型在管道的启动时间和频率方面有所不同:

  • 连续:管道在不停止的情况下运行。 仅当次秒延迟是硬性要求时使用,因为持续计算成本更高。 与文件事件配对。

  • 计划:管道按时间表运行(例如,每小时运行一次)。 当延迟要求宽松(分钟到小时)时使用。 可配合目录列表工作,但即使在定时模式下,文件事件也能通过避免全量目录扫描来降低成本。

有关使用 Trigger.AvailableNow 进行批处理调度的详细信息,请参阅 使用 Trigger.AvailableNow 和速率限制

选择正确的文件发现模式

自动加载程序支持三种文件发现模式,在设置复杂性、可伸缩性和成本方面具有不同的权衡。

模式 配置复杂性 Scalability Cost 何时使用
文件事件(推荐) 较低(只需一次权限设置) 每小时数百万个文件 Lowest 大多数工作负荷的默认值
经典文件通知 高(21+ 个云配置选项) 每小时数百万个文件 中等 文件事件不可用时
目录列表 没有 受目录大小限制 最高(LIST API 成本) 小型目录、一次性回填或安全策略阻止文件事件时

文件事件通过为每个外部位置使用一个订阅和一个队列,而不是为每个流分别使用一个订阅和一个队列,来整合云存储资源。 性能差异非常大规模:目录列表必须扫描每个触发器上的整个源目录,因此引入时间随目录大小而增加。 文件事件直接传送新的文件通知,因此,无论目录中有多少对象,引入时间都保持低。

启用文件事件

文件事件需要一次性云权限授权,以及已配置为使用托管文件事件服务的外部位置。 设置完成后,所有从该外部位置读取的 Auto Loader 流都可以使用文件事件,而无需额外配置。

  1. 在云提供商端授予所需的云权限。 要求因云提供商而异。 请参阅 为外部位置设置文件事件

  2. 在您的 Auto Loader 查询中,将 cloudFiles.useManagedFileEvents 设置为 true

    df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.useManagedFileEvents", "true")
      .load("/path/to/data/dir"))
    

当无法使用文件事件时

在以下情况下,可能无法使用文件事件:

  • 外部位置未配置文件事件。
  • 组织安全策略不允许在共享外部位置上启用文件事件。

在这些情况下,请使用 经典文件通知模式目录列表模式。 有关文件检测模式的完整比较,请参阅 “比较自动加载程序文件检测模式”。

管理架构演变

自动加载程序自动推断架构,但配置架构演变的方式会影响数据完整性和管道稳定性。 使用下表选择策略。

Scenario Recommendation
架构已知且已修复 提供显式架构 .schema()
架构未知,需要累加更改 schemaEvolutionMode: addNewColumns
架构未知,预计类型会发生变化 schemaEvolutionMode: addNewColumnsWithTypeWidening
必须遵守严格的架构协定 schemaEvolutionMode: failOnNewColumns
任意或不可预知的架构 作为 Variant 类型引入

选择策略后,请应用以下做法来微调架构演变的行为方式。

对已知字段类型使用架构提示

使用此选项 cloudFiles.schemaHints 对事先知道的字段强制实施类型,同时仍允许对其他字段进行架构推理。

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaHints", "id long, amount double")
  .load("/path/to/data/dir"))

对于兼容的类型更改,使用类型拓宽

addNewColumnsWithTypeWidening架构演变模式会自动扩大兼容类型(例如int,为long),而不是将数据路由到_rescued_data列。 这样就无需通过后处理作业来处理简单的类型提升。

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "parquet")
  .option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
  .load("/path/to/data/dir"))

Variant 类型引入不可预测的架构

如果您的数据不符合任何特定架构,或者架构不断变化,请将数据作为 Variant 类型摄取。

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("singleVariantColumn", "data")
  .load("/path/to/data/dir"))

Variant 在查询时提供架构读取,但效率低于查询结构化列。 有关架构推理和演变的完整机制,请参阅 自动加载器中的配置架构推理和演变

处理不良数据和数据质量

以下做法有助于在将数据传播到下游层之前检测、捕获和隔离不良数据。

启用 _rescued_data_corrupt_record

Auto Loader 提供了两列,用于捕获无法被正常解析的数据。

  • _rescued_data 捕获与当前架构不匹配的字段。 它由 Auto Loader 自动添加。
  • _corrupt_record 记录完全无法解析的行。 使用 columnNameOfCorruptRecord 启用它:
df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaHints", "_corrupt_record string")
  .option("columnNameOfCorruptRecord", "_corrupt_record")
  .load("/path/to/data/dir"))

Databricks 建议使用 columnNameOfCorruptRecord 而不是 badRecordsPath,以避免可能导致遗漏损坏记录的竞争条件。

使用 Lakeflow 管道预期规则进行监控

设置 Lakeflow 管道预期,以验证 _rescued_data_corrupt_record 在正常情况下为 NULL。 非 NULL 值表示架构偏移或数据损坏。

import dlt

@dlt.table
@dlt.expect("no rescued data", "_rescued_data IS NULL")
@dlt.expect("no corrupt records", "_corrupt_record IS NULL")
def bronze_table():
    return (spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.schemaHints", "_corrupt_record string")
        .option("columnNameOfCorruptRecord", "_corrupt_record")
        .load("/path/to/data/dir"))

隔离损坏的数据

将包含无法解析数据的行隔离到专用接收端中,以便调查。 这可以防止损坏的数据传播到下游层。

import dlt

@dlt.table
def corrupt_records_sink():
    return dlt.read_stream("bronze_table").where("_corrupt_record IS NOT NULL")

@dlt.view
def clean_table():
    return dlt.read_stream("bronze_table").where("_corrupt_record IS NULL")

使用源文件元数据批注数据

在 Auto Loader 摄取查询中包含 _metadata 列。 至少捕获 file_pathfile_modification_time。 这样,您就可以将数据问题追溯到特定的源文件,并与 cloud_files_state() 联接,以获取完整的文件生命周期信息。

df = (spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .load("/path/to/data/dir")
  .select("*", "_metadata.file_path", "_metadata.file_modification_time"))

有关详细信息,请参阅 文件元数据列

优化成本和性能

以下做法减少了自动加载程序三个主要成本驱动因素:云 LIST API 调用、空闲计算和长期存储增长。

  • 使用文件事件将 LIST API 成本降到最低:文件事件提供增量文件发现,无需每次运行的完整目录列表。 这是对 Auto Loader 而言成效最大的成本优化措施。

  • 将文件到达触发器用于事件驱动的处理:文件到达触发器仅在新文件到达时启动管道,因此无需为空闲计算付费。

  • 使用 cloudFiles.cleanSource 存档已处理的文件:用于 cloudFiles.cleanSource 自动删除或移动已处理的文件。 这可降低长期流的存储成本和目录列表成本。 有关完整详细信息,请参阅 源目录中的存档文件以降低成本

    • 使用 delete 模式在引入后删除文件。
    • 使用 move 模式将文件存档到其他位置进行合规性或审核。
    df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.cleanSource", "delete")
      .load("/path/to/data/dir"))
    

    Warning

    如果多个自动加载程序流或其他客户端从同一源目录读取,请不要启用 cloudFiles.cleanSource

  • 利用性能改进:升级到最新的 Databricks Runtime 或使用无服务器计算从最近的自动加载程序性能改进中受益。

检查点管理

检查点存储流的进度和文件状态。 错误配置或丢失检查点需要完全重启,因此请将其视为关键基础结构。

  • 切勿将云对象生命周期策略应用于检查点位置。 如果删除检查点文件,则流状态已损坏,并且必须从头开始重启。
  • 为每个流和源目录使用单独的检查点。
  • 对于生命周期长且流量大的流,可考虑使用 cloudFiles.maxFileAge 以限制状态增长。 采用保守设置(最低建议为 90 天)。 如果将此值设置得过高,一旦文件超出该时间窗口,可能会导致自动加载程序已导入的文件被重复处理。

有关完整详细信息,请参阅 文件事件跟踪

利用卷实现基于文件事件的优化文件发现

为了提高文件事件的处理性能,请为自动加载程序读取的每个路径或子目录创建一个外部卷。 向 Auto Loader 提供卷路径(例如,/Volumes/catalog/schema/volume),而不是云路径(例如,s3://bucket/path)。 这会通过优化的数据访问模式优化文件发现。