Spark 1.6(Python 2.7)下读写压缩SequenceFile:求可用测试代码
嘿,我知道你在CDH 5.12 Quickstart虚拟机上用Spark 1.6 + Python 2.7折腾带压缩的SequenceFile读写时踩了坑——之前找到的Scala代码直接套Python肯定不行,毕竟俩API的写法差不少。下面给你一套能直接跑通的测试代码,还附上关键细节说明,帮你快速搞定:
可直接运行的测试代码
1. 写入带压缩的SequenceFile
from pyspark import SparkContext # 初始化SparkContext sc = SparkContext(appName="CompressedSequenceFileWriter") # 构造简单的测试键值对数据 test_rdd = sc.parallelize([(1, "spark"), (2, "hadoop"), (3, "cdh"), (4, "python")]) # 指定压缩编码器的完整Java类名 # CDH 5.12内置支持的编码器可选: # - Gzip: org.apache.hadoop.io.compress.GzipCodec # - Snappy: org.apache.hadoop.io.compress.SnappyCodec compression_codec = "org.apache.hadoop.io.compress.GzipCodec" # 保存到HDFS路径(如果要写本地文件,前缀加file://,比如file:///home/cloudera/test_seq) test_rdd.saveAsSequenceFile( path="/user/hadoop/compressed_seq_test", compressionCodecClass=compression_codec ) # 关闭SparkContext sc.stop()
2. 读取带压缩的SequenceFile
from pyspark import SparkContext # 初始化SparkContext sc = SparkContext(appName="CompressedSequenceFileReader") # 读取SequenceFile,指定键值对应的Hadoop Writable类 # Spark会自动识别文件的压缩格式,无需手动指定编码器 read_rdd = sc.sequenceFile( path="/user/hadoop/compressed_seq_test", keyClass="org.apache.hadoop.io.IntWritable", valueClass="org.apache.hadoop.io.Text" ) # 转换Writable对象为Python原生类型并打印结果 print("读取到的SequenceFile内容:") for key_writable, value_writable in read_rdd.collect(): # 数值类型Writable调用get()获取值,Text类型调用toString() key = key_writable.get() value = value_writable.toString() print(f"Key: {key}, Value: {value}") # 关闭SparkContext sc.stop()
关键注意事项
- 压缩编码器的正确写法:必须使用完整的Java类名,不能只写类名。如果需要用LZO压缩,得先在CDH集群上安装LZO支持,对应的类名是
com.hadoop.compression.lzo.LzoCodec。 - Writable对象转换:读取出来的是Hadoop的Writable对象,不是Python原生类型,必须调用对应的方法(比如
get()、toString())才能拿到实际值。 - 路径选择:默认路径是HDFS路径,如果你要操作本地文件,记得加上
file://前缀,比如file:///home/cloudera/local_seq_test。 - 版本适配:这套代码完全适配Spark 1.6和Python 2.7,和Spark 2.x及以上版本的API写法会有差异,别搞混了。
内容的提问来源于stack exchange,提问作者singhak.bhu
相关产品推荐
相关产品推荐

