如何配置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点自动执行脚本(需替换为你的Python路径和脚本路径)0 */12 * * * /usr/bin/python3 /path/to/your/spark_script.py - 云环境/大数据调度工具:如Apache Airflow、Azure Data Factory、AWS CloudWatch Events等,可根据你的部署环境配置周期性任务
内容的提问来源于stack exchange,提问作者MMV
相关产品推荐
相关产品推荐

