Python中DataFrame写入文件无报错但文件为空问题求助
问题排查与修复建议
核心错误点
- 错误调用DataFrame的
write方法而非文件对象的write
处理hive_df部分时,你错误使用了hive_df.write()——这是PySpark DataFrame用于将数据写入外部存储(如HDFS、数据库)的API,并非写入本地文件的方法。正确操作应该是调用output_file.write(),写入之前生成的hive_df_string。 - 缩进语法错误
if not os.path.exists(local_filename):后的open(local_filename, 'a').close()没有缩进,这会触发Python语法错误;同时try块、except以及debug打印的缩进混乱,会导致代码逻辑无法正常执行。 - 手动关闭文件对象多余且易出问题
使用with open(...) as output_file:时,with语句会自动在代码块结束时关闭文件,无需手动调用output_file.close(),手动关闭可能导致文件提前关闭,后续写入操作失效。 - 换行符转义错误
代码中"Data from Oracle "这类写法不符合Python语法,字符串内的换行需要用\n,比如"Data from Oracle\n",否则会触发语法错误。
修复后的代码示例
import datetime import os current_date = datetime.datetime.now().strftime("%Y%m%d") local_filename = local_path + file_name + "_" + current_date + ".csv" try: # 'w'模式会自动创建不存在的文件,无需提前创建 with open(local_filename, 'w') as output_file: output_file.write("Data from Oracle\n") output_file.write(','.join(oracle_df.columns) + '\n') oracle_df_string = oracle_df.rdd.map(lambda row: ','.join(map(str, row))).collect() output_file.write('\n'.join(oracle_df_string) + '\n') output_file.write("==============================\n") output_file.write("Data from Hive\n") output_file.write(','.join(hive_df.columns) + '\n') hive_df_string = hive_df.rdd.map(lambda row: ','.join(map(str, row))).collect() output_file.write('\n'.join(hive_df_string) + '\n') output_file.write("==============================\n") print("Data successfully written to file :", local_filename) except IOError as (errno, strerror): print("I/O error({0}): {1}".format(errno, strerror)) # Debug prints print("Oracle DataFrame:") oracle_df.show() print("oracle_df.columns values", oracle_df.columns) print("oracle_df_string", oracle_df_string)
额外优化建议
- 避免提前创建文件:使用
'w'模式打开文件时,若文件不存在会自动创建,无需提前用open(local_filename, 'a').close()操作。 - 使用PySpark内置CSV写入方法:如果仅需将DataFrame写入CSV文件,推荐使用PySpark原生API,更高效且避免手动处理RDD与字符串拼接的繁琐,示例:
该方法适配大数据量场景,无需手动处理列名与行数据的拼接。# 写入Oracle数据到CSV oracle_df.write.mode("overwrite").option("header", "true").csv(local_filename.replace('.csv', '_oracle')) # 写入Hive数据到CSV hive_df.write.mode("overwrite").option("header", "true").csv(local_filename.replace('.csv', '_hive'))
内容的提问来源于stack exchange,提问作者Rahul Patidar
相关产品推荐
相关产品推荐

