Pyspark:如何以非DataFrame Schema的变量名作为目录写入Parquet文件?
在PySpark中使用Schema外变量作为Parquet写入目录的解决方法
你的代码报错是因为误解了df.write.parquet()的参数用法:第二个参数会被Spark识别为分区列名,而table_name并不是DataFrame Schema里的字段,所以才会提示它不属于Schema。
要实现用循环变量作为目录名写入Parquet,正确的做法是直接把变量拼接到输出路径中:
import os tables_list = ['abc','def','xyz'] for table_name in tables_list: # 将table_name拼接到根路径后,作为每个循环的输出目录 output_path = os.path.join("s3://bucket/output/", table_name) df.write.parquet(output_path)
如果需要覆盖已存在的目录,可以添加mode参数:
df.write.mode("overwrite").parquet(output_path)
补充说明:partitionBy()方法才是用于按DataFrame内的字段分区存储的,这要求字段必须存在于Schema中;而你只是想用外部变量作为目录名,直接拼接路径即可,不需要涉及分区参数。
内容的提问来源于stack exchange,提问作者Manish
相关产品推荐
相关产品推荐

