使用Rust的Polars库在Azure Data Lake读写CSV转Parquet遇问题
使用Rust Polars读写Azure Data Lake文件遇到的问题
我尝试使用Rust的polars库从Azure Data Lake读取CSV文件,随后将其以Parquet文件形式写回Azure Data Lake。Polars官网声称支持读写所有常见文件及云存储,包括Azure Storage:
Polars supports reading and writing to all common files (e.g. csv, json, parquet), cloud storage (S3, Azure Blob, BigQuery) and databases (e.g. postgres, mysql).
然而相关示例极少且均与AWS相关,我无法找到读取CSV的示例,现有示例均为读取Parquet文件。
我尝试将Parquet文件写入Azure Storage,使用URI abfss://data@<account_name>.dfs.core.windows.net/output/bronze/customers.parquet时出现如下错误:
thread 'main' panicked at 'called `Result::unwrap()` on an `Err` value: ComputeError(ErrString("unknown url scheme"))'
Rust代码实现
use clap::Parser; use cloud::AzureConfigKey; use polars::prelude::*; /* How to run: cargo run -- --input abfss://data@<account_name>.dfs.core.windows.net/raw/customers.csv --output abfss://data@<account_name>.dfs.core.windows.net/output/bronze/customers.parquet */ /// Simple program to convet CSV to Parquet #[derive(Parser, Debug)] #[command(author, version, about, long_about = None)] struct Args { /// Input file path #[arg(short, long, value_name = "INPUT-PATH")] input: String, /// Output file path #[arg(short, long, value_name = "OUTPUT-PATH")] output: String, } fn main() -> PolarsResult<()> { dotenvy::dotenv().unwrap(); let args = Args::parse(); let input_path = args.input; let output_path = args.output; let account_name = dotenvy::var("ACCOUNT_NAME").unwrap(); let access_key = dotenvy::var("ACCESS_KEY").unwrap(); let sas_key = dotenvy::var("SAS_KEY").unwrap(); // Propagate the credentials and other cloud options. let cloud_options = cloud::CloudOptions::default().with_azure([ (AzureConfigKey::AccountName, account_name.to_string()), (AzureConfigKey::AccessKey, access_key.to_string()), (AzureConfigKey::ContainerName, "data".to_string()), (AzureConfigKey::SasKey, sas_key.to_string()) ]); let mut schema = Schema::new(); schema.with_column("CustomerID".to_string().into(), DataType::Int64); schema.with_column("FirstName".to_string().into(), DataType::Utf8); schema.with_column("LastName".to_string().into(), DataType::Utf8); schema.with_column("FullName".to_string().into(), DataType::Utf8); // FIXME: There was no documention about how to read csv from cloud let df_csv = CsvReader::from_path(input_path)? .has_header(true) .with_dtypes(Some(Arc::new(schema))) .with_encoding(CsvEncoding::LossyUtf8) .finish()?; println!("{}", df_csv); let cloud_options = Some(cloud_options); // FIXME: This function returns an error that unknown url scheme df_csv.lazy() .sink_parquet_cloud( output_path, cloud_options, Default::default(), ) .unwrap(); Ok(()) }
Cargo.toml配置
[package] name = "poolstar" version = "0.1.0" edition = "2021" [dependencies] dotenvy = "0.15.7" chrono = "0.4.31" clap = { version = "4.4.4" , features = ["derive"] } polars = { version = "0.33.2", features = ["lazy", "temporal", "describe", "json", "parquet", "dtype-datetime", "polars-io", "cloud", "azure", "cloud_write"] }
内容的提问来源于stack exchange,提问作者DSaad
相关产品推荐
相关产品推荐

