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

使用Rusoto实现CSV上传S3并导入Redshift时遇编译报错求助

问题解决:Rusoto上传S3到Redshift的代码错误修复

错误原因分析

  • CopyCommand导入错误:Rusoto Redshift库中不存在CopyCommand类型,COPY是Redshift的SQL命令,无需导入Rusoto结构体。
  • ExecuteStatementMessage和execute_statement方法不存在:Rusoto的RedshiftClient仅用于管理Redshift集群资源(如创建集群、修改参数),不支持执行SQL语句。执行Redshift SQL需使用PostgreSQL兼容客户端,比如tokio-postgres。

修正后的代码

首先更新Cargo.toml依赖:

[dependencies]
rusoto_core = "0.48.0"
rusoto_s3 = "0.48.0"
rusoto_credential = "0.48.0"
config = "0.13"
tokio = { version = "1.0", features = ["full"] }
tokio-postgres = "0.7"
postgres-native-tls = "0.5"
native-tls = "0.2"

修改后的源代码:

extern crate rusoto_core;
extern crate rusoto_s3;
extern crate rusoto_credential;
extern crate config;

use rusoto_core::Region;
use rusoto_s3::S3;
use rusoto_s3::{S3Client, PutObjectRequest};
use rusoto_credential::DefaultCredentialsProvider;
use std::fs::File;
use std::io::Read;
use config::Config;
use tokio_postgres::{Client, NoTls};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 加载配置文件
    let mut settings = Config::default();
    settings.merge(config::File::with_name("config.toml"))?;

    // 读取S3配置
    let s3_bucket = settings.get_str("aws.s3_bucket")?;
    let s3_key = settings.get_str("aws.s3_key")?;

    // 读取Redshift配置
    let redshift_username = settings.get_str("redshift.username")?;
    let redshift_password = settings.get_str("redshift.password")?;
    let redshift_host = settings.get_str("redshift.host")?;
    let redshift_port = settings.get_str("redshift.port")?;
    let redshift_db = settings.get_str("redshift.database").unwrap_or("dwh");

    // 读取本地CSV文件
    let mut csv_file = File::open("exchange_rates.csv")?;
    let mut csv_data = Vec::new();
    csv_file.read_to_end(&mut csv_data)?;

    // 创建S3客户端并上传文件
    let credentials_provider = DefaultCredentialsProvider::new()?;
    let s3_client = S3Client::new_with(
        rusoto_core::HttpClient::new().expect("Failed to create HTTP client"),
        credentials_provider,
        Region::default(),
    );

    let s3_put_request = PutObjectRequest {
        bucket: s3_bucket.clone(),
        key: s3_key.clone(),
        body: Some(csv_data.into()),
        ..Default::default()
    };
    let response = s3_client.put_object(s3_put_request).await?;
    println!("S3 upload successful. ETag: {:?}", response.e_tag);

    // 通过PostgreSQL协议连接Redshift
    let conn_str = format!(
        "postgresql://{}:{}@{}:{}/{}",
        redshift_username, redshift_password, redshift_host, redshift_port, redshift_db
    );
    let (client, connection) = tokio_postgres::connect(&conn_str, NoTls).await?;

    // 后台维护连接
    tokio::spawn(async move {
        if let Err(e) = connection.await {
            eprintln!("Redshift connection error: {}", e);
        }
    });

    // 构造COPY命令(推荐用IAM角色替代硬编码密钥)
    let copy_command = format!(
        "COPY bi_etl.erbe__rust FROM 's3://{}/{}' \
        IAM_ROLE 'arn:aws:iam::YOUR_ACCOUNT_ID:role/YourRedshiftS3AccessRole' \
        DELIMITER ',' CSV;",
        s3_bucket, s3_key
    );

    // 执行COPY命令
    client.execute(&copy_command, &[]).await?;
    println!("COPY command executed successfully");

    Ok(())
}

关键修改说明

  1. 移除无用的rusoto_redshift依赖及错误导入,因为Rusoto Redshift客户端不负责SQL执行。
  2. 引入tokio-postgres客户端,通过PostgreSQL协议连接Redshift执行COPY命令。
  3. 替换凭证方式:推荐使用IAM角色(需提前为Redshift集群关联具备S3读取权限的IAM角色),避免硬编码密钥,提升安全性。若坚持使用密钥,可换回原格式,但不推荐。
  4. 添加Redshift连接的后台任务,确保连接持续有效。

额外注意事项

  • 确认Redshift安全组允许你的机器访问默认端口5439。
  • 确认S3桶权限配置正确:IAM角色需具备S3读权限;若用密钥,密钥对应用户需有S3读权限。
  • 确认Redshift目标表bi_etl.erbe__rust已存在,且结构与CSV文件匹配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 03:33:27