Sruhvx-jpg commented on code in PR #3165: URL: https://github.com/apache/iceberg-rust/pull/3165#discussion_r4236606727
########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! `object_store`-based storage implementation for Apache Iceberg. +//! +//! This crate provides [`ObjectStoreStorage`] and [`ObjectStoreStorageFactory`], +//! which implement the [`Storage`] and +//! [`StorageFactory`] traits from the `iceberg` crate +//! using the [`object_store`](https://docs.rs/object_store) crate as the backend. +//! +//! Currently only S3 storage is supported (via the `object_store-s3` feature flag, +//! enabled by default). + +#[cfg(feature = "object_store-s3")] +mod s3; + +use std::ops::Range; +use std::sync::Arc; + +use async_trait::async_trait; +use bytes::{Bytes, BytesMut}; +use dashmap::DashMap; +use futures::stream::BoxStream; +use futures::{StreamExt, TryStreamExt}; +#[cfg(feature = "object_store-s3")] +use iceberg::io::S3Config; +use iceberg::io::{ + FileMetadata, FileRead, FileWrite, InputFile, OutputFile, Storage, StorageConfig, + StorageFactory, +}; +use iceberg::{Error, ErrorKind, Result}; +use object_store::path::Path as ObjectStorePath; +use object_store::{MultipartUpload, ObjectStore, ObjectStoreExt, PutPayload}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// Convert an `object_store::Error` into an `iceberg::Error`, +/// dispatching known variants to their corresponding `ErrorKind`. +fn from_object_store_error(e: object_store::Error) -> Error { + let (kind, msg) = match &e { + object_store::Error::NotFound { path, .. } => { + (ErrorKind::Unexpected, format!("Object not found: {path}")) + } + object_store::Error::AlreadyExists { path, .. } => ( + ErrorKind::Unexpected, + format!("Object already exists: {path}"), + ), + object_store::Error::PermissionDenied { path, .. } => { + (ErrorKind::Unexpected, format!("Permission denied: {path}")) + } + object_store::Error::Unauthenticated { path, .. } => { + (ErrorKind::Unexpected, format!("Unauthenticated: {path}")) + } + object_store::Error::NotSupported { .. } => ( + ErrorKind::FeatureUnsupported, + "Operation not supported".to_string(), + ), + _ => ( + ErrorKind::Unexpected, + "Failure in doing io operation".to_string(), + ), + }; + Error::new(kind, msg).with_source(e) +} + +/// Property key for configuring S3 bulk delete batch size (matches Iceberg Java specification). +pub const S3_DELETE_BATCH_SIZE: &str = "s3.delete.batch-size"; +/// Default batch size for S3 bulk deletions (matches Iceberg Java specification). +pub const DEFAULT_DELETE_BATCH_SIZE: usize = 250; +/// Maximum batch size allowed by the AWS S3 DeleteObjects API specification. +pub const S3_MAX_DELETE_BATCH_SIZE: usize = 1000; + +fn parse_delete_batch_size(config: &StorageConfig) -> usize { + if let Some(val) = config.get(S3_DELETE_BATCH_SIZE) { + match val.parse::<usize>() { + Ok(parsed) if parsed > 0 => parsed.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + _ => { + tracing::warn!( + val = %val, + "Invalid s3.delete.batch-size; falling back to default 250" + ); + DEFAULT_DELETE_BATCH_SIZE + } + } + } else { + DEFAULT_DELETE_BATCH_SIZE + } +} + +fn default_delete_batch_size() -> usize { + DEFAULT_DELETE_BATCH_SIZE +} + +/// Convert `object_store::ObjectMeta` into `iceberg::io::FileMetadata`. +fn to_file_metadata(meta: object_store::ObjectMeta) -> FileMetadata { + FileMetadata { size: meta.size } +} + +/// `object_store`-based storage factory. +/// +/// Use this factory with `FileIOBuilder::new(factory)` to create FileIO instances +/// backed by the `object_store` crate. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum ObjectStoreStorageFactory { + /// S3 storage factory. + #[cfg(feature = "object_store-s3")] + S3, +} + +#[typetag::serde(name = "ObjectStoreStorageFactory")] +impl StorageFactory for ObjectStoreStorageFactory { + #[allow(unused_variables)] + fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>> { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorageFactory::S3 => { + let s3_config = S3Config::try_from(config)?; + let delete_batch_size = parse_delete_batch_size(config); + tracing::info!( + batch_size = delete_batch_size, + "Initialized S3 storage with delete batch size {} (configure via '{}' to adjust)", + delete_batch_size, + S3_DELETE_BATCH_SIZE + ); + Ok(Arc::new(ObjectStoreStorage::S3(S3Storage { + config: Arc::new(s3_config), + delete_batch_size, + store_cache: Arc::new(DashMap::new()), + }))) + } + } + } +} + +type StoreCache = Arc<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[serde(default = "default_delete_batch_size")] + delete_batch_size: usize, + #[serde(skip, default)] + store_cache: StoreCache, +} + +/// `object_store`-based storage implementation. +/// +/// Stores are cached per bucket to avoid rebuilding the client on every operation. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum ObjectStoreStorage { + /// S3 storage variant. + #[cfg(feature = "object_store-s3")] + S3(S3Storage), +} + +struct StoreAndPath { + bucket: String, + store: Arc<dyn ObjectStore>, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc<dyn ObjectStore>, + locations: Vec<ObjectStorePath>, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result<StoreAndPath> { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => { + let parsed = parse_s3_url(path)?; + + let store = s3 + .store_cache + .entry(parsed.bucket.clone()) + .or_try_insert_with(|| build_s3_store(&s3.config, &parsed.bucket))? + .value() + .clone(); + + let object_path = ObjectStorePath::from(parsed.relative.as_str()); Review Comment: leaving this un-resolved till u approve of the updated test -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
