使用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(©_command, &[]).await?; println!("COPY command executed successfully"); Ok(()) }
关键修改说明
- 移除无用的
rusoto_redshift依赖及错误导入,因为Rusoto Redshift客户端不负责SQL执行。 - 引入
tokio-postgres客户端,通过PostgreSQL协议连接Redshift执行COPY命令。 - 替换凭证方式:推荐使用IAM角色(需提前为Redshift集群关联具备S3读取权限的IAM角色),避免硬编码密钥,提升安全性。若坚持使用密钥,可换回原格式,但不推荐。
- 添加Redshift连接的后台任务,确保连接持续有效。
额外注意事项
- 确认Redshift安全组允许你的机器访问默认端口5439。
- 确认S3桶权限配置正确:IAM角色需具备S3读权限;若用密钥,密钥对应用户需有S3读权限。
- 确认Redshift目标表
bi_etl.erbe__rust已存在,且结构与CSV文件匹配。
内容的提问来源于stack exchange,提问作者Luc
相关产品推荐
相关产品推荐

