Python调用Spark newAPIHadoopFile修改记录分隔符遇错求助
我来帮你搞定这个问题!先理清之前代码失效和报错的原因,再给你正确的实现方案。
为什么sc.textFile不生效?
sc.textFile()底层默认用的是旧版Hadoop MapReduce API(org.apache.hadoop.mapred.TextInputFormat),而你设置的textinputformat.record.delimiter是给新版API(org.apache.hadoop.mapreduce.lib.input.TextInputFormat)用的参数,两者不匹配,所以配置根本没被读取到,自然不会生效。
为什么newAPIHadoopFile报错?
你碰到的AttributeError: 'SparkConf' object has no attribute '_get_object_id',是因为newAPIHadoopFile的最后一个参数需要的是Hadoop原生的Configuration对象,不是Spark的SparkConf。在PySpark里,我们得通过SparkContext的hadoopConfiguration属性来设置Hadoop相关配置,而不是直接传SparkConf。
正确的实现代码
下面是修正后的newAPIHadoopFile方案,这也是推荐的做法:
appName = "My Test" fname = "myfile.txt" from pyspark import SparkContext, SparkConf if __name__ == "__main__": conf = SparkConf().setAppName(appName) sc = SparkContext(conf=conf) sc.setLogLevel("ERROR") # 关键步骤:通过sc.hadoopConfiguration设置Hadoop的分隔符参数 sc.hadoopConfiguration.set("textinputformat.record.delimiter", "H") # 调用newAPIHadoopFile,参数对应新版Hadoop API的类 rdd2 = sc.newAPIHadoopFile( fname, "org.apache.hadoop.mapreduce.lib.input.TextInputFormat", "org.apache.hadoop.io.LongWritable", "org.apache.hadoop.io.Text" ) # 提取实际的文本记录(返回的元组第一个元素是偏移量,第二个是内容) result = rdd2.map(lambda x: x[1]).collect() print("improved:", result)
备选方案:用textFile配合旧API参数
如果你更习惯用sc.textFile,可以设置旧版API对应的分隔符参数mapred.textinputformat.record.delimiter,代码如下:
appName = "My Test" fname = "myfile.txt" from pyspark import SparkContext, SparkConf if __name__ == "__main__": conf = SparkConf().setAppName(appName) # 设置旧版API的分隔符参数 conf.set("mapred.textinputformat.record.delimiter", "H") sc = SparkContext(conf=conf) sc.setLogLevel("ERROR") rdd1 = sc.textFile(fname) print("normal:", rdd1.collect())
注意:旧版API(mapred)在较新的Hadoop/Spark版本中已被标记为过时,所以优先推荐第一种newAPIHadoopFile的方案。
内容的提问来源于stack exchange,提问作者vy32

