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
相关产品推荐
相关产品推荐

