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

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机制,完全可控。

步骤与代码示例

  1. 添加依赖(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>
    
  2. 核心读取代码(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读取。

步骤与代码示例

  1. 添加依赖:
    同方案一,需包含zstd-jni依赖。

  2. 自定义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; }
        }
    }
    
  3. 注册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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 10:14:59