You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 08:28:26