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

PySpark遍历S_ID列表统计各表新增行数并匹配Snowflake存量计数的实现方法

修正实现方案

原代码问题梳理

  • 过滤条件df.S_ID == 'i'将i硬编码为字符串常量,未使用循环变量的实际取值
  • 循环内直接覆盖循环变量i的值,导致原S_ID标识丢失
  • Snowflake查询语句中S_ID = 'i'同样为硬编码,无法匹配不同的S_ID取值
  • 使用列表存储结果无法满足{表名:新增行数}的键值对格式要求,需改用字典存储

修正后代码

# 初始化字典存储最终统计结果
new_counts_dict = {}

for s_id in mylist:
    # 统计当前S_ID对应的上传数据总行数
    upload_total = df.filter(df.S_ID == s_id).count()
    # 拼接动态查询语句,代入当前S_ID取值
    sf_query = f"SELECT R_ID FROM mytable WHERE S_ID = '{s_id}'"
    # 统计Snowflake中当前表的已有行数
    existing_total = spark.read.format(SNOWFLAKE_SOURCE_NAME)\
                                .options(**sfOptions)\
                                .option("query", sf_query)\
                                .load()\
                                .count()
    # 计算新增行数并转字符串格式
    new_rows = str(upload_total - existing_total)
    # 写入结果字典,key为S_ID对应表名,value为新增行数
    new_counts_dict[s_id] = new_rows

结果使用说明

运行完成后new_counts_dict即为符合要求的{表名:新增行数}格式数据,可直接遍历输出或做后续业务处理,输出示例代码如下:

for table_name, add_count in new_counts_dict.items():
    print(f"{table_name}表新增行数:{add_count}")

注意:如果你的S_ID取值存在引号等特殊字符,建议使用Snowflake参数化查询替代字符串拼接,避免SQL注入风险。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 16:45:04