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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 19:42:44