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

如何配置PySpark EventHub流作业运行10分钟后停止并定时执行?

修改PySpark流作业实现定时停止与周期性执行

一、让流作业运行10分钟后自动停止

修改原代码,利用StreamingQuery.awaitTermination(timeout)的超时参数控制运行时长,超时后主动停止流作业:

# 设置运行时长为10分钟(转换为毫秒)
RUN_DURATION_MS = 10 * 60 * 1000

streamer = (
    spark.readStream.format("eventhubs")
    .options(**ehConf)
    .load()
    .writeStream.foreachBatch(write_to_parquet_table)
    .option(
        "checkpointLocation",
        eventhub_checkpoint_location,
    )
    .outputMode("update")
    .start()
)

try:
    # 等待指定时长,超时后返回False
    if not streamer.awaitTermination(RUN_DURATION_MS):
        print("流作业已运行10分钟,开始停止")
        streamer.stop()
        # 等待作业完全终止
        streamer.awaitTermination()
except Exception as e:
    print(f"流作业运行异常: {str(e)}")
    streamer.stop()
    raise

关键说明:

  • awaitTermination(RUN_DURATION_MS)会阻塞直到流作业终止或达到超时时间,超时后返回False
  • 调用streamer.stop()主动终止流作业,之后再次调用awaitTermination()确保作业完全停止
  • 异常捕获块保证出现错误时也能正确停止流,避免资源泄漏

二、实现每12小时执行一次脚本

PySpark本身不适合内置定时调度,推荐使用外部工具触发脚本:

  • Linux/macOS cron:执行crontab -e编辑定时任务,添加:
    0 */12 * * * /usr/bin/python3 /path/to/your/spark_script.py
    
    该规则会在每天0点、12点自动执行脚本(需替换为你的Python路径和脚本路径)
  • 云环境/大数据调度工具:如Apache Airflow、Azure Data Factory、AWS CloudWatch Events等,可根据你的部署环境配置周期性任务

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 12:45:17