PySpark循环中print输出异常,current_date值不匹配的问题排查与解决
这个问题我之前处理过,本质是Python的GIL(全局解释器锁)在IO操作时释放,导致循环变量被提前更新,加上PySpark的操作涉及集群IO,触发了线程调度切换,才出现这种随机的日期不匹配情况。
问题原因分析
你代码里的current_date是循环的迭代变量,在Python中,当执行涉及IO的操作(比如Spark读取Parquet的元数据阶段),GIL会被暂时释放,这时候Python的线程调度器可能会切换到循环的下一次迭代,把current_date更新为下一个日期。当IO操作完成、GIL重新获取后,回到当前迭代的代码执行print(f'loop continued: {current_date}')时,current_date已经是下一次迭代的值了——这就是为什么两次print的日期会不一样,而且问题随机出现(取决于IO操作的耗时和线程调度时机)。
另外,你可能会疑惑:spark.read.parquet不是只创建DataFrame吗?为什么会有IO?其实Spark在创建DataFrame时会先读取元数据(比如分区信息),这一步涉及IO操作,足以触发GIL释放。
解决方法
针对这个问题,有几种简单有效的解决方式,你可以根据自己的场景选择:
1. 固定当前迭代的日期引用
在循环内部把current_date赋值给一个局部变量,确保后续操作使用的是当前迭代的固定值,不受循环变量更新影响:
from datetime import timedelta for current_date in dates_list: # 固定当前日期的引用,避免后续被循环更新 current_date_fixed = current_date print(f'loop started: {current_date_fixed}') loop_start_date = current_date_fixed - timedelta(days=90) dates_ = [loop_start_date + timedelta(days=i) for i in range((current_date_fixed - loop_start_date).days + 1)] dates = ','.join(map(lambda date: date.strftime("%Y-%m-%d"), dates_)) res = spark.read.parquet(f'my_path/day={{{dates}}}') print(f'loop continued: {current_date_fixed}') print('----------')
因为datetime对象是不可变的,直接赋值就可以固定引用,不用担心被修改。
2. 触发同步Action,阻塞直到Spark操作完成
在spark.read.parquet之后,执行一个同步的Spark Action(比如count()),确保当前迭代的Spark操作完全完成后,再进入下一次循环:
from datetime import timedelta for current_date in dates_list: print(f'loop started: {current_date}') loop_start_date = current_date - timedelta(days=90) dates_ = [loop_start_date + timedelta(days=i) for i in range((current_date - loop_start_date).days + 1)] dates = ','.join(map(lambda date: date.strftime("%Y-%m-%d"), dates_)) res = spark.read.parquet(f'my_path/day={{{dates}}}') # 触发Action,阻塞当前线程直到读取完成 res.count() print(f'loop continued: {current_date}') print('----------')
这样Spark操作会同步执行,GIL不会提前释放,循环变量也就不会被提前更新。
3. 把迭代逻辑封装成函数
将每个日期的处理逻辑封装成独立函数,把current_date作为参数传入——函数内部的参数是局部变量,和循环的迭代变量完全隔离,不会受后续迭代影响:
from datetime import timedelta def process_single_date(current_date): print(f'loop started: {current_date}') loop_start_date = current_date - timedelta(days=90) dates_ = [loop_start_date + timedelta(days=i) for i in range((current_date - loop_start_date).days + 1)] dates = ','.join(map(lambda date: date.strftime("%Y-%m-%d"), dates_)) res = spark.read.parquet(f'my_path/day={{{dates}}}') print(f'loop continued: {current_date}') print('----------') for current_date in dates_list: process_single_date(current_date)
这种方式代码更清晰,也从根本上避免了变量引用冲突的问题。
备注:内容来源于stack exchange,提问作者mohammad hassan bigdeli shamlo

