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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 02:25:52