Python的csv.writer(open())写法如何转换为PySpark可运行代码
可行解决方案
直接将不等长的每行数据提前拼接为|分隔的单个字符串,再通过PySpark的文本写入接口输出,完全匹配原生Python的输出逻辑,同时规避原生open函数不兼容的问题。
方案1:小数据量(数据存储在Driver端内存)实现
from operator import itemgetter # 原始数据&排序逻辑和原有Python逻辑完全一致 my_list = [ [3, 'ab','ac','ad'], [4, 'ae','af','at','aj','aa'], [1, 'ar','aa','at','as'], [2, 'ay','au','aa','ar','aa','a1'] ] sorted_list = sorted(my_list, key=itemgetter(0)) # 每行转为|分隔的字符串,直接构造单列Spark DataFrame写入 row_str_list = ['|'.join(map(str, row)) for row in sorted_list] spark.createDataFrame(row_str_list, "string") \ .coalesce(1) \ # 需单文件输出时加这行,不需要可删除 .write \ .mode("overwrite") \ .text("file:///本地输出目录路径 或 HDFS/对象存储路径")
方案2:大数据量全分布式处理(数据已存储为RDD,无需加载到Driver端)
from operator import itemgetter # 示例为将本地列表转为RDD,实际场景可直接读取数据源生成RDD my_list = [ [3, 'ab','ac','ad'], [4, 'ae','af','at','aj','aa'], [1, 'ar','aa','at','as'], [2, 'ay','au','aa','ar','aa','a1'] ] rdd = spark.sparkContext.parallelize(my_list) # 分布式排序+拼接字符串,直接写入存储 sorted_rdd = rdd.sortBy(lambda x: x[0]) \ .map(lambda row: '|'.join(map(str, row))) sorted_rdd.coalesce(1) \ # 需单文件输出时加这行,不需要可删除 .saveAsTextFile("file:///本地输出目录路径 或 HDFS/对象存储路径")
原方案失败原因
- 不等长列表转为Pandas DataFrame时会自动补全NaN空值,Spark会以最长行的长度定义表结构,写入时会将补全的空值一并输出,和原有Python代码的输出逻辑不符
- 若数据量较大,将全量数据加载到Driver端转Pandas会存在内存溢出风险
注意事项
- Spark写入的是输出目录而非单个文件,加
coalesce(1)后输出目录下的part-00000文件即为和原生Python写入结果完全一致的目标文件 - 写入本地路径时路径前缀加
file:///,写入HDFS/对象存储直接填对应路径即可,无需调用原生open函数
内容的提问来源于stack exchange,提问作者Tanai Goncalves
相关产品推荐
相关产品推荐

