循环提取指定时间段数据至多个DataFrame的技术求助
问题描述
我有4年的员工离职数据,需要提取指定日期前1个月的时间段数据,并分别存储到不同的DataFrame中:例如month_1=2023-04-25对应df_1,存储2023-03-25至2023-04-25的数据;month_2=2023-03-25对应df_2,存储2023-02-25至2023-03-25的数据。尝试用for循环执行Spark SQL查询时,循环能正常运行,但无法将结果保存到不同DataFrame;修改代码后出现错误,寻求代码修正方案。
初始尝试代码:
#months of interest date_col = month_1, month_2 for d in date_col: date = datetime.datetime.strptime(d, '%Y-%m-%d') starttime = (date + relativedelta(months=-1)).strftime("%Y-%m-%d") endtime = (date + relativedelta(months=0)).strftime("%Y-%m-%d") print('STARTTIME' + ' '+ starttime) print('ENDTIME' + ' '+ endtime) Current_month_resign_emp = spark.sql(""" SELECT * FROM Table_resignation_emp WHERE record_update >= '{0}' AND record_update <= '{1}' """.format(starttime, endtime))
修改后报错的代码:
#months of interest date_col = month_1, month_2 result_month = {} for d in date_col: date = datetime.datetime.strptime(d, '%Y-%m-%d') starttime = (date + relativedelta(months=-1)).strftime("%Y-%m-%d") endtime = (date + relativedelta(months=0)).strftime("%Y-%m-%d") print('STARTTIME' + ' '+ starttime) print('ENDTIME' + ' '+ endtime) Current_month_resign_emp = spark.sql(""" SELECT * FROM Table_resignation_emp WHERE record_update >= '{0}' AND record_update <= '{1}' """.format(starttime, endtime)) result_month.append(Current_month_resign_emp) Current_month_resign_emp[d] = pd.DataFrame(Current_month_resign_emp) print(result_month)
报错原因:字典对象无append方法,且错误地将DataFrame赋值给自身的键。
错误分析
- 字典使用错误:
result_month是字典类型,字典没有append()方法,应该通过键值对赋值来添加元素。 - DataFrame转换与赋值错误:
Current_month_resign_emp[d] = pd.DataFrame(Current_month_resign_emp)逻辑混乱,试图给Spark DataFrame对象添加键值对;同时Spark DataFrame不能直接用pd.DataFrame()构造,需调用toPandas()方法完成转换。
修正后的代码
import datetime from dateutil.relativedelta import relativedelta # 定义目标日期(示例值) month_1 = "2023-04-25" month_2 = "2023-03-25" date_col = (month_1, month_2) # 用字典存储每个时间段的DataFrame,方便按标识或日期访问 result_dfs = {} for idx, target_date in enumerate(date_col, start=1): # 解析日期并计算起止时间 date_obj = datetime.datetime.strptime(target_date, '%Y-%m-%d') start_time = (date_obj + relativedelta(months=-1)).strftime("%Y-%m-%d") end_time = date_obj.strftime("%Y-%m-%d") # 简化写法,等价于原relativedelta(months=0) print(f'STARTTIME {start_time}') print(f'ENDTIME {end_time}') # 执行Spark SQL查询获取数据 spark_df = spark.sql(""" SELECT * FROM Table_resignation_emp WHERE record_update >= '{0}' AND record_update <= '{1}' """.format(start_time, end_time)) # 将Spark DataFrame转换为Pandas DataFrame(如果不需要Pandas,可直接存储spark_df) pd_df = spark_df.toPandas() # 存储到字典,用df_1、df_2作为键,也可以用原日期字符串当键 result_dfs[f'df_{idx}'] = pd_df # 单独提取每个DataFrame df_1 = result_dfs['df_1'] df_2 = result_dfs['df_2'] print("已生成的DataFrame标识:", result_dfs.keys())
关键修正点
- 用字典的键值对赋值替代错误的
append()方法,通过result_dfs[f'df_{idx}'] = pd_df保存每个时间段的DataFrame。 - 使用Spark DataFrame的
toPandas()方法正确转换为Pandas DataFrame,避免直接用pd.DataFrame()构造的错误。 - 用
enumerate生成df_1、df_2这类标识,方便后续直接访问指定的DataFrame;也可以直接用目标日期字符串作为字典的键。
内容的提问来源于stack exchange,提问作者Ichichaa
相关产品推荐
相关产品推荐

