You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Databricks Autoloader与现有Workflow Job集成以自动触发S3文件处理

集成Databricks Autoloader Notebook与现有Workflow Job实现S3 Zip文件自动化处理

这里提供三种可行的集成方案,根据你的批量处理需求和实时性要求选择:

方案1:Autoloader Notebook直接调用Jobs API触发现有任务

适合需要单个文件实时触发处理的场景,在Autoloader查询到新增文件后,直接调用Databricks Jobs API启动现有Job并传递文件路径参数。

实现步骤:

  • 在Autoloader Notebook中,完成新增文件查询后,添加API调用逻辑:
    import requests
    import json
    
    # 从Secret Scope读取API Token(提前在Workspace中配置好)
    token = dbutils.secrets.get(scope="your-secret-scope", key="databricks-api-token")
    workspace_url = "https://<你的Databricks工作区URL>"
    target_job_id = "<现有Workflow Job的ID>" # 在Job详情页可查看
    
    # 从cloud_files_state查询获取新增文件路径列表
    checkpoint_query = "SELECT * FROM cloud_files_state('%s') ORDER BY create_time DESC" % (checkpoint_path)
    new_files_df = spark.sql(checkpoint_query)
    new_files = [row.path for row in new_files_df.collect()]
    
    # 遍历文件触发Job
    for file_path in new_files:
        payload = {
            "job_id": target_job_id,
            "notebook_params": {
                "s3_file_path": file_path # 与现有Job中配置的参数名保持一致
            }
        }
        response = requests.post(
            f"{workspace_url}/api/2.1/jobs/run-now",
            headers={"Authorization": f"Bearer {token}"},
            json=payload
        )
        if response.ok:
            print(f"已触发Job处理文件: {file_path}")
        else:
            print(f"触发失败,错误: {response.text}")
    
  • 注意事项:确保运行Notebook的服务主体有调用Jobs API的权限,Secret Scope已正确配置API Token。

方案2:将Autoloader作为Workflow的前置任务串联执行

适合批量处理场景,把Autoloader Notebook加入现有Workflow作为第一个任务,通过任务间参数传递把新增文件列表传给后续处理任务。

实现步骤:

  1. 修改Autoloader Notebook,在查询到新增文件后用dbutils.notebook.exit()返回参数:
    import json
    
    # 查询新增文件逻辑
    checkpoint_query = "SELECT * FROM cloud_files_state('%s') ORDER BY create_time DESC" % (checkpoint_path)
    new_files_df = spark.sql(checkpoint_query)
    new_files = [row.path for row in new_files_df.collect()]
    
    # 返回文件列表(转成JSON字符串方便传递)
    dbutils.notebook.exit(json.dumps(new_files))
    
  2. 编辑现有Workflow Job:
    • 添加新的Notebook任务,选择你的Autoloader Notebook作为第一个任务
    • 对后续的处理任务,开启任务输入,选择接收前置Autoloader任务的输出
  3. 在处理任务的Notebook中接收参数并处理:
    import json
    
    # 获取前置任务传递的文件列表
    files_input = dbutils.widgets.get("task_input")
    new_files = json.loads(files_input)
    
    # 遍历文件执行处理逻辑(或适配现有任务的参数逻辑)
    for file_path in new_files:
        # 这里调用你的处理代码,或设置现有任务的s3_file_path参数
        dbutils.widgets.text("s3_file_path", file_path)
        # 执行原有处理流程
    
  • 进阶:如果需要逐个文件处理,可以使用Workflow的任务循环功能,把Autoloader返回的文件列表作为循环迭代的输入,每个循环实例处理单个文件。

方案3:基于Delta表触发Job(批量延迟处理)

适合非实时的批量处理场景,让Autoloader把新增文件路径写入Delta表,再配置Job监听Delta表更新自动触发。

实现步骤:

  1. 调整Autoloader逻辑,将新增文件路径写入Delta表:
    (spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "binaryFile")
        .option("pathGlobFilter", "*.zip") # 只监听zip文件
        .load("s3://你的存储桶路径/")
        .select("path", "modificationTime")
        .writeStream
        .format("delta")
        .option("checkpointLocation", checkpoint_path)
        .trigger(availableNow=True) # 按需触发,或用processingTime设置间隔
        .table("s3_zip_files_log")
    )
    
  2. 编辑现有Workflow Job:
    • 进入Job的触发器设置,选择Delta表更新
    • 选择刚才创建的s3_zip_files_log表,配置触发条件(如每次有新数据写入时触发)
  3. 在处理任务中读取Delta表的新增数据进行批量处理:
    # 读取Delta表中的新增文件(可通过版本或时间范围过滤)
    new_files_df = spark.read.table("s3_zip_files_log").filter("modificationTime > current_timestamp() - interval 1 hour")
    # 执行批量处理逻辑
    

内容的提问来源于stack exchange,提问作者Saravanan Ponnaiah

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 01:55:18