PySpark中如何将换行分隔符改为?配置后仍用 问题求助
解决PySpark自定义记录分隔符不生效的问题
你的问题出在两点:一是Hadoop输入格式的配置没正确传递,二要确保特殊分隔符的表示正确。直接用SparkConf设置textinputformat.record.delimiter不会被textFile正确识别,因为这是Hadoop层面的配置,需要通过SparkContext的Hadoop配置对象来设置。
另外注意:你要替换的分隔符是ASCII控制字符SOH(十六进制0x01),在Python里得用\x01表示;如果是代码里写的,那对应\x02,要和文件里的实际分隔符严格对应。
修改后的代码方案1(推荐,适配textFile)
from pyspark import SparkContext, SparkConf conf = SparkConf().setAppName('example').setMaster('local[*]') sc = SparkContext(conf=conf) # 关键:通过SparkContext获取Hadoop配置并设置分隔符 sc._jsc.hadoopConfiguration().set("textinputformat.record.delimiter", "\x01") # 现在读取文件会使用指定的分隔符 rdd = sc.textFile("MY_PATH")
方案2(更底层的newAPIHadoopFile,兼容性更强)
如果方案1还是不生效,可以直接调用Hadoop API读取,更明确控制输入格式:
from pyspark import SparkContext, SparkConf from org.apache.hadoop.mapreduce.lib.input import TextInputFormat from org.apache.hadoop.io import Text, LongWritable conf = SparkConf().setAppName('example').setMaster('local[*]') sc = SparkContext(conf=conf) hadoop_conf = sc._jsc.hadoopConfiguration() hadoop_conf.set("textinputformat.record.delimiter", "\x01") # 读取后返回(key, value)对,取value作为记录内容 rdd = sc.newAPIHadoopFile( "MY_PATH", TextInputFormat, LongWritable, Text, conf=hadoop_conf ).map(lambda item: item[1].toString())
原理说明:textinputformat.record.delimiter是Hadoop TextInputFormat的专属配置项,必须设置到Hadoop的配置实例中才会生效,仅通过SparkConf传递无法被Hadoop的输入格式读取到。SparkContext初始化后,通过sc._jsc.hadoopConfiguration()拿到的才是Hadoop的实际配置对象,修改它才能让文件读取时使用自定义分隔符。
内容的提问来源于stack exchange,提问作者dotan
相关产品推荐
相关产品推荐

