如何配置Databricks工作流,使通知任务按动态获取时间启动
Databricks工作流配置:按动态获取的时间启动通知任务
问题描述
我有一个Databricks工作流,计划完成两项操作:
- 运行Notebook获取当日日落的UTC时间;
- 在该时间点通过短信向自己发送通知。
用于获取日落时间的代码如下:
import requests from datetime import datetime url = 'https://api.openweathermap.org/data/2.5/weather?lat={lat}&lon={lon}&appid={appid}' response = requests.get(url) weather = response.json() sunset_utc = weather["sys"]["sunset"] sunset_utc_formatted = datetime.fromtimestamp(sunset_utc).strftime('%Y-%m-%d %H:%M:%S') dbutils.jobs.taskValues.set(key = "Time", value = sunset_utc_formatted) sunset_utc_formatted
这段代码可以生成格式化后的日落时间(如2024-08-12 01:17:51),且短信通知Notebook单独测试正常,但原工作流设置会导致通知任务运行失败,需要配置工作流让“Notification”任务按照“GetTime”任务获取的特定时间启动。
解决方案
Databricks原生工作流的任务依赖默认是任务完成后立即触发下一个任务,无法直接基于动态生成的时间延迟触发,可通过以下两种方式实现需求:
方式1:在GetTime任务中调用API调度Notification任务
这是最可靠的方案,适合长时间延迟的场景:
- 先在工作流中创建好可独立运行的Notification任务,测试确认正常后记录该任务的ID;
- 修改GetTime任务的代码,在获取日落时间后,调用Databricks Jobs API调度指定任务在目标时间启动:
import requests from datetime import datetime import json # 原获取日落时间的逻辑 url = 'https://api.openweathermap.org/data/2.5/weather?lat={lat}&lon={lon}&appid={appid}' response = requests.get(url) weather = response.json() sunset_utc = weather["sys"]["sunset"] sunset_utc_formatted = datetime.fromtimestamp(sunset_utc).strftime('%Y-%m-%d %H:%M:%S') dbutils.jobs.taskValues.set(key = "Time", value = sunset_utc_formatted) # 新增:调度Notification任务 databricks_host = dbutils.notebook.entry_point.getDbutils().notebook().getContext().apiUrl().get() databricks_token = dbutils.notebook.entry_point.getDbutils().notebook().getContext().apiToken().get() notification_job_id = "替换为你的Notification任务ID" schedule_url = f"{databricks_host}/api/2.1/jobs/run-now" headers = {"Authorization": f"Bearer {databricks_token}", "Content-Type": "application/json"} payload = { "job_id": notification_job_id, "start_time": sunset_utc_formatted # 指定UTC时间启动任务 } # 发送调度请求并处理结果 response = requests.post(schedule_url, headers=headers, data=json.dumps(payload)) response.raise_for_status()
方式2:在Notification任务中添加延迟逻辑
适合短时间延迟的场景(几小时内),但会占用集群资源:
- 在GetTime任务中计算当前时间到日落时间的秒数差,并传递给下一个任务:
# 原获取日落时间代码... current_utc = datetime.utcnow().timestamp() delay_seconds = sunset_utc - current_utc dbutils.jobs.taskValues.set(key = "DelaySeconds", value = delay_seconds)
- 修改Notification任务的代码,先等待指定时长再执行通知:
import time # 获取传递的延迟时间 delay_seconds = dbutils.jobs.taskValues.get(taskKey="GetTime", key="DelaySeconds") # 等待到目标时间 time.sleep(delay_seconds) # 执行短信通知的原有代码
关键注意事项
- 方式1中,需确保GetTime任务所在集群有调用Databricks Jobs API的权限;
start_time参数必须是UTC格式字符串,与代码生成的sunset_utc_formatted格式一致;- 若使用方式1,建议移除工作流中GetTime与Notification任务的直接依赖,避免重复触发。
内容的提问来源于stack exchange,提问作者MyNameHere
相关产品推荐
相关产品推荐

