Spark 2.4.7+Hadoop2.7.7环境下如何读取Zstandard压缩Parquet文件
解决方案:在Spark 2.4.7+Hadoop 2.7.7环境读取S3上Zstd压缩的Parquet文件
针对你遇到的环境限制(无法安装系统组件、只能修改作业代码),有两种可行的方案实现读取Zstd压缩的Parquet文件为Spark DataFrame:
方案一:手动读取+解压缩+Parquet解析
通过AWS SDK读取S3上的文件二进制流,用纯Java的Zstd库解压缩,再借助Parquet Java API解析成Spark Row,最终生成DataFrame。这种方案不依赖Hadoop的Codec机制,完全可控。
步骤与代码示例
添加依赖(Maven为例):
打包作业时需包含这些纯Java依赖,无需系统级安装:<dependency> <groupId>com.github.luben</groupId> <artifactId>zstd-jni</artifactId> <version>1.5.5-11</version> <!-- 兼容Spark 2.4.7的Java 8环境 --> </dependency> <dependency> <groupId>org.apache.parquet</groupId> <artifactId>parquet-hadoop</artifactId> <version>1.10.1</version> <!-- Spark 2.4.7内置的Parquet版本 --> </dependency> <dependency> <groupId>software.amazon.awssdk</groupId> <artifactId>s3</artifactId> <version>2.20.100</version> <!-- 或使用旧版AWS SDK v1 --> </dependency>核心读取代码(Java实现,可转Scala):
import com.github.luben.zstd.ZstdInputStream; import org.apache.parquet.hadoop.ParquetReader; import org.apache.parquet.hadoop.api.GenericReadSupport; import org.apache.parquet.hadoop.util.HadoopInputFile; import org.apache.parquet.schema.MessageType; import org.apache.parquet.schema.SchemaParser; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.RowFactory; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.types.StructType; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.ListObjectsV2Request; import software.amazon.awssdk.services.s3.model.S3Object; import java.io.InputStream; import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; public class ZstdParquetReader { public static Dataset<Row> readFromS3(SparkSession spark, String s3Path, String parquetSchemaStr) { // 1. 解析Parquet Schema并转换为Spark Schema MessageType parquetSchema = new SchemaParser().parseMessageType(parquetSchemaStr); StructType sparkSchema = org.apache.spark.sql.parquet.ParquetSchemaConverter.convert(parquetSchema); // 2. 拆分S3路径为桶名和前缀,列出所有Parquet文件 String[] pathParts = s3Path.split("/", 4); String bucket = pathParts[2]; String prefix = pathParts.length > 3 ? pathParts[3] : ""; S3Client s3Client = S3Client.create(); List<String> objectKeys = s3Client.listObjectsV2(ListObjectsV2Request.builder() .bucket(bucket) .prefix(prefix) .build()) .contents().stream() .map(obj -> obj.key()) .filter(key -> key.endsWith(".parquet")) .collect(Collectors.toList()); // 3. 并行读取文件并转换为Spark Row return spark.createDataFrame( spark.sparkContext().parallelize(objectKeys) .toJavaRDD() .flatMap(key -> { List<Row> rows = new ArrayList<>(); try (S3Object s3Obj = s3Client.getObject(b -> b.bucket(bucket).key(key)); InputStream rawStream = s3Obj.content(); ZstdInputStream zstdStream = new ZstdInputStream(rawStream); ParquetReader<Object> reader = ParquetReader.builder(new GenericReadSupport(), HadoopInputFile.fromStream(zstdStream, new org.apache.hadoop.fs.Path(s3Path))) .withConf(spark.sparkContext().hadoopConfiguration()) .build()) { Object record; while ((record = reader.read()) != null) { org.apache.parquet.generic.GenericRecord genericRecord = (org.apache.parquet.generic.GenericRecord) record; Object[] values = parquetSchema.getFields().stream() .map(field -> genericRecord.get(field.getName())) .toArray(); rows.add(RowFactory.create(values)); } } catch (Exception e) { throw new RuntimeException("Failed to process file: " + key, e); } return rows.iterator(); }), sparkSchema ); } }
方案二:自定义Hadoop CompressionCodec
实现Hadoop的CompressionCodec接口,内部用纯Java的Zstd库处理解压缩,注册到Spark配置后,直接使用原生spark.read.parquet读取。
步骤与代码示例
添加依赖:
同方案一,需包含zstd-jni依赖。自定义Zstd Codec实现:
import com.github.luben.zstd.ZstdInputStream; import com.github.luben.zstd.ZstdOutputStream; import org.apache.hadoop.io.compress.*; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; public class CustomZstdCodec implements CompressionCodec { @Override public CompressionInputStream createInputStream(InputStream in) throws IOException { return new ZstdCompressionInputStream(in); } @Override public CompressionInputStream createInputStream(InputStream in, CompressionDecompressor decompressor) throws IOException { return createInputStream(in); } @Override public CompressionOutputStream createOutputStream(OutputStream out) throws IOException { return new ZstdCompressionOutputStream(out); } @Override public CompressionOutputStream createOutputStream(OutputStream out, CompressionCompressor compressor) throws IOException { return createOutputStream(out); } @Override public Class<? extends CompressionDecompressor> getDecompressorType() { return DummyDecompressor.class; } @Override public Class<? extends CompressionCompressor> getCompressorType() { return DummyCompressor.class; } @Override public CompressionDecompressor createDecompressor() { return new DummyDecompressor(); } @Override public CompressionCompressor createCompressor() { return new DummyCompressor(); } @Override public String getDefaultExtension() { return ".zst"; } // 包装ZstdInputStream的自定义CompressionInputStream static class ZstdCompressionInputStream extends CompressionInputStream { private final ZstdInputStream zstdIn; public ZstdCompressionInputStream(InputStream in) throws IOException { super(in); this.zstdIn = new ZstdInputStream(in); } @Override public int read() throws IOException { return zstdIn.read(); } @Override public int read(byte[] b, int off, int len) throws IOException { return zstdIn.read(b, off, len); } @Override public void resetState() throws IOException {} } // 包装ZstdOutputStream的自定义CompressionOutputStream static class ZstdCompressionOutputStream extends CompressionOutputStream { private final ZstdOutputStream zstdOut; public ZstdCompressionOutputStream(OutputStream out) throws IOException { super(out); this.zstdOut = new ZstdOutputStream(out); } @Override public void write(int b) throws IOException { zstdOut.write(b); } @Override public void write(byte[] b, int off, int len) throws IOException { zstdOut.write(b, off, len); } @Override public void finish() throws IOException { zstdOut.finish(); } } // 空实现的Decompressor(仅满足接口要求,实际解压缩由ZstdInputStream处理) static class DummyDecompressor implements CompressionDecompressor { @Override public int decompress(byte[] b, int off, int len) throws IOException { return 0; } @Override public boolean needsInput() { return false; } @Override public void setInput(byte[] b, int off, int len) {} @Override public void reset() {} @Override public void end() {} @Override public int getRemaining() { return 0; } } // 空实现的Compressor(仅用于满足接口,无需压缩功能可忽略) static class DummyCompressor implements CompressionCompressor { @Override public int compress(byte[] b, int off, int len) throws IOException { return 0; } @Override public boolean needsInput() { return false; } @Override public void setInput(byte[] b, int off, int len) {} @Override public void reset() {} @Override public void end() {} @Override public float getProgress() { return 0; } } }注册Codec并读取文件(Scala示例):
import org.apache.spark.sql.SparkSession object ZstdParquetJob { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("ZstdParquetReader") // 注册自定义Codec .config("spark.hadoop.io.compression.codecs", "com.yourpackage.CustomZstdCodec") .getOrCreate() // 直接使用原生Parquet读取API val df = spark.read.parquet("s3://your-bucket/path/to/zstd-parquet/") df.show() spark.stop() } }
关键注意事项
- 打包依赖:所有第三方依赖(zstd-jni、AWS SDK等)必须打包到作业的fat jar中,确保集群节点能加载到这些类。
- Schema处理:建议手动指定Parquet Schema,避免动态推断带来的性能开销和潜在错误;可通过本地解压缩一个样本文件获取Schema字符串。
- 资源调整:针对大文件场景,调整Spark执行器的内存和核心数,避免并行读取时出现OOM。
- 兼容性测试:自定义Codec方案需测试Spark Parquet读取逻辑是否能正确识别并调用自定义Codec,部分复杂Parquet结构可能需要额外适配。
内容的提问来源于stack exchange,提问作者Nikhil Pandit
相关产品推荐
相关产品推荐

