如何在Azure Databricks中运行永久/长期流处理工作流并实现故障重启?
Azure Event Hub流处理任务持久化与调度配置问题
我们需要监听Azure Event Hub并将数据写入Azure Databricks的Delta表,已编写如下流处理代码:
df = spark.readStream.format("eventhubs").options(**ehConf).load() # Code omitted where message content is expanded into columns in the dataframe df.writeStream \ .format("delta") \ .outputMode("append") \ .option("checkpointLocation", "/tmp/delta/events/_checkpoints/") \ .toTable("mydb.mytable")
该代码运行正常,但Notebook会停留在df.writeStream行直到任务被取消。请问应如何配置才能让该工作流永久运行,且在代码崩溃时自动重启?是否可以将其设为常规工作流,设置为每分钟运行且最大并发数为1?
解决方案
一、让流处理任务永久运行并自动重启
- 用Databricks作业部署流处理:不要在交互式Notebook中运行流处理,把代码封装成Notebook或Python脚本,创建Databricks作业:
- 选择合适的集群:推荐使用自动缩放集群,或固定集群并开启自动重启机制
- 任务配置:选择对应任务类型(Notebook/脚本)关联你的代码
- 开启自动重试:在任务的高级选项中,设置重试规则(比如无限重试、指定间隔重试),确保任务崩溃后自动重启
- 修复检查点路径:当前使用的
/tmp/属于临时存储,集群重启后会丢失,必须改为DBFS持久化路径,比如dbfs:/delta/events/_checkpoints/,这样任务重启后能从上次中断的位置继续消费,避免数据丢失或重复处理
二、改成每分钟运行的常规批处理工作流
完全可行,但需要调整代码逻辑,将流式读取改为批处理模式:
- 调整读取逻辑:把
readStream替换为read,通过时间范围或偏移量控制每次读取的数据,示例代码如下:from datetime import datetime, timedelta # 读取过去1分钟的Event Hub数据 end_time = datetime.utcnow() start_time = end_time - timedelta(minutes=1) ehConf["startingTimestamp"] = start_time.isoformat() ehConf["endingTimestamp"] = end_time.isoformat() df = spark.read.format("eventhubs").options(**ehConf).load() # 消息解析列的逻辑保持不变 df.write.format("delta").mode("append").saveAsTable("mydb.mytable") - 配置作业调度:创建Databricks作业,设置调度频率为每分钟,同时在作业的并发控制中设置最大并发数为1,确保上一批处理完成后再启动下一批,避免Delta表写入冲突
- 偏移量管理:如果要精准避免重复读取,建议将每次读取的结束偏移量存储到一个小表中,下次读取时从该偏移量开始,比时间戳更可靠
内容的提问来源于stack exchange,提问作者Mathias Rönnlund
相关产品推荐
相关产品推荐

