如何在PySpark中循环执行代码直至满足指定条件
补全PySpark脚本的循环与迭代提示功能
以下是修改完成的完整脚本,已补充循环逻辑和每5次迭代的提示功能:
import time from pyspark.sql import SparkSession import pyspark.sql.functions as f # Spark session builder spark = SparkSession.builder.getOrCreate() count_max = 50 def oracle_pull_check(): df = spark.read \ .format("jdbc") \ .option("url", "jdbc:oracle:thin:@your_aliastns?TNS_ADMIN=path/to/wallet") \ .option("dbtable", 'table_name or query') \ .option("user", "user") \ .option("password", "password") \ .option("driver", "oracle.jdbc.driver.OracleDriver") \ .load() # 创建count_limit列,判断count_now是否大于等于count_max df1 = df.withColumn("count_max", f.lit(count_max))\ .withColumn("count_limit", f.when(f.col("count_now") >= f.col("count_max"), 'Y').otherwise(f.lit('N'))) # 统计未满足条件的记录数 return df1.filter(f.col("count_limit") == 'N').count() iteration_count = 0 while True: limit_values = oracle_pull_check() if limit_values == 0: print("所有记录均已满足count_limit条件") break iteration_count += 1 # 每5次迭代(即每5分钟)打印提示 if iteration_count % 5 == 0: print(f"检查已运行{iteration_count}分钟") # 等待1分钟后重试 time.sleep(60)
关键修改说明:
- 循环逻辑:使用
while True实现持续检查,直至所有记录满足条件后通过break退出循环 - 迭代计数:新增
iteration_count变量跟踪重试次数,每次未满足条件时递增 - 定时提示:通过
iteration_count % 5 == 0判断每5次迭代(对应5分钟),打印运行时长提示 - 退出逻辑:当
limit_values == 0(所有记录满足条件)时,打印完成信息并终止循环
内容的提问来源于stack exchange,提问作者nmr
相关产品推荐
相关产品推荐

