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

EC2环境下Spark读取Redshift数据时S3临时目录读取过慢求助

Spark从Redshift读取数据后读取S3临时目录速度过慢的问题与优化方案

我在EC2实例上使用Apache Spark从AWS Redshift读取数据,再将结果以CSV格式写入AWS S3存储桶。当前使用io.github.spark_redshift_community.spark.redshift驱动,该驱动会先执行查询并将结果存入S3临时目录的CSV文件中。由于特定约束,无法使用Athena或UNLOAD命令。功能已实现,但从S3临时目录读取数据的过程异常缓慢:仅1万条共2MB的数据,读取并写入目标S3位置耗时近一分钟。日志显示Redshift写入临时目录速度较快,延迟完全出现在读取临时目录阶段。运行Spark的EC2实例已拥有S3存储桶的IAM角色访问权限。

读取Redshift的代码

spark.read()
            .format("io.github.spark_redshift_community.spark.redshift")
            .option("url",URL)
            .option("query",    QUERY)
            .option("user", USER_ID)
            .option("password", PASSWORD)
            .option("tempdir", TEMP_DIR)
            .option("forward_spark_s3_credentials", "true")
            .load();

pom.xml依赖配置

<dependencies>

    <dependency>
        <groupId>com.eclipsesource.minimal-json</groupId>
        <artifactId>minimal-json</artifactId>
        <version>0.9.5</version>
    </dependency>
    <dependency>
        <groupId>org.apache.maven.plugins</groupId>
        <artifactId>maven-assembly-plugin</artifactId>
        <version>3.3.0</version>
    </dependency>
    <dependency>
        <groupId>org.ini4j</groupId>
        <artifactId>ini4j</artifactId>
        <version>0.5.4</version>
    </dependency>
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <version>1.18.2</version>
    </dependency>
    <dependency>
        <groupId>org.slf4j</groupId>
        <artifactId>slf4j-api</artifactId>
        <version>1.7.26</version>
    </dependency>

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-avro_2.12</artifactId>
        <version>3.3.1</version>
    </dependency>

    <dependency>
        <groupId>io.github.spark-redshift-community</groupId>
        <artifactId>spark-redshift_2.12</artifactId>
        <version>4.2.0</version>
    </dependency>

    <dependency>
        <groupId>io.delta</groupId>
        <artifactId>delta-core_2.12</artifactId>
        <version>2.2.0</version>
    </dependency>

    <dependency>
        <groupId>org.scala-lang</groupId>
        <artifactId>scala-library</artifactId>
        <version>2.12.15</version>
    </dependency>

    <dependency>
        <groupId>org.apache.hadoop</groupId>
        <artifactId>hadoop-aws</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>com.amazonaws</groupId>
        <artifactId>aws-java-sdk-s3</artifactId>
        <version>1.12.389</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>com.amazonaws</groupId>
        <artifactId>aws-java-sdk-bundle</artifactId>
        <version>1.12.389</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-hadoop-cloud_2.12</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_2.12</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>org.apache.hadoop</groupId>
        <artifactId>hadoop-common</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>org.apache.hadoop</groupId>
        <artifactId>hadoop-client</artifactId>
        <version>3.3.1</version>
        <scope>provided</scope>
    </dependency>

    <dependency>
        <groupId>junit</groupId>
        <artifactId>junit</artifactId>
        <version>3.8.1</version>
        <scope>test</scope>
    </dependency>
</dependencies>

优化方案

1. 调整Spark读取S3的配置参数

修改Spark的Hadoop配置,提升S3客户端并发与缓存能力:

spark.sparkContext().hadoopConfiguration().set("fs.s3a.connection.maximum", "100");
spark.sparkContext().hadoopConfiguration().set("fs.s3a.fast.upload", "true");
spark.sparkContext().hadoopConfiguration().set("fs.s3a.buffer.dir", "/tmp");
spark.sparkContext().hadoopConfiguration().set("fs.s3a.metadata.cache.enable", "true");
  • fs.s3a.connection.maximum:提升S3连接池最大数量,增加并发读取能力
  • fs.s3a.fast.upload:启用快速上传模式,优化小文件读写效率
  • fs.s3a.buffer.dir:指定EC2本地磁盘作为缓冲区,减少网络IO开销
  • fs.s3a.metadata.cache.enable:缓存S3文件元数据,避免重复请求

2. 减少临时目录的小文件数量

Redshift写入临时目录时生成的大量小文件会增加Spark读取的开销,可通过以下方式优化:

  • 在Redshift查询中调整结果分区,比如通过ORDER BY或DISTINCT让Redshift生成更少的输出文件
  • 添加redshift.jdbc.fetchsize参数(如10000),控制JDBC读取批次大小,间接减少临时文件数量

3. 优化EC2与S3的网络环境

  • 确保EC2实例与S3临时桶处于同一AWS区域,避免跨区域网络延迟
  • 若使用VPC部署,配置S3网关端点,让EC2直接通过内网访问S3,绕过公网
  • 检查EC2实例的网络带宽占用情况,排除其他进程抢占带宽的可能

4. 升级依赖版本

  • 尝试将spark-redshift_2.12升级至最新稳定版,修复已知的S3读取性能问题
  • 确保hadoop-aws与aws-java-sdk-s3版本和Spark版本兼容,避免依赖冲突导致的性能损耗

5. 优化临时目录位置

将临时目录设置为目标S3桶的子目录,减少跨桶操作的额外开销,避免跨桶数据传输的延迟。

内容的提问来源于stack exchange,提问作者Ishan Sanganeria

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:55:18