From 3c0e2a8d553c4e7d1dbf5fe16227f8f991dca6ed Mon Sep 17 00:00:00 2001 From: jeremyhi Date: Mon, 28 Sep 2026 15:22:02 +0000 Subject: [PATCH] feat: export Metric snapshots with packed Parquet objects (#9382) * feat: export Metric snapshots with packed Parquet objects Signed-off-by: jeremyhi * fix: reject empty packed export time ranges Signed-off-by: jeremyhi * test: cover cancellation during packed export I/O Signed-off-by: jeremyhi --------- Signed-off-by: jeremyhi --- Cargo.lock | 2 + src/cli/Cargo.toml | 1 + src/cli/src/data/export_v2.rs | 4 +- src/cli/src/data/export_v2/command.rs | 69 +- src/cli/src/data/export_v2/coordinator.rs | 190 ++++- src/cli/src/data/export_v2/data.rs | 4 + src/cli/src/data/export_v2/manifest.rs | 2 +- src/cli/src/data/snapshot_storage.rs | 102 ++- src/cli/src/database.rs | 102 +-- src/common/datasource/Cargo.toml | 1 + src/common/datasource/src/lib.rs | 1 + src/common/datasource/src/packed_writer.rs | 669 ++++++++++++++++++ src/common/datasource/src/parquet_writer.rs | 116 ++- src/operator/src/statement/copy_table_to.rs | 54 +- src/operator/src/statement/database_copy.rs | 2 +- src/operator/src/statement/export_database.rs | 70 +- .../src/statement/export_logical_tables.rs | 30 +- .../statement/export_logical_tables/tests.rs | 39 +- .../export_logical_tables/writers.rs | 9 +- src/servers/src/http.rs | 6 +- .../tests/export_logical_tables.rs | 196 ++++- .../cases/snapshot_parquet_copy/setup.sql | 6 +- .../cases/snapshot_parquet_copy/verify.result | 33 +- .../cases/snapshot_parquet_copy/verify.sql | 3 + 24 files changed, 1584 insertions(+), 127 deletions(-) create mode 100644 src/common/datasource/src/packed_writer.rs diff --git a/Cargo.lock b/Cargo.lock index 4ced42fc860..96903faaa25 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2088,6 +2088,7 @@ dependencies = [ "table", "tempfile", "tokio", + "tokio-util", "url", "uuid", ] @@ -2397,6 +2398,7 @@ dependencies = [ "paste", "regex", "serde", + "serde_json", "snafu 0.8.6", "strum 0.27.1", "tokio", diff --git a/src/cli/Cargo.toml b/src/cli/Cargo.toml index 8178fa1b0b5..6abe504a3ab 100644 --- a/src/cli/Cargo.toml +++ b/src/cli/Cargo.toml @@ -71,3 +71,4 @@ common-test-util.workspace = true object-store = { workspace = true, features = ["testing"] } serde.workspace = true tempfile.workspace = true +tokio-util.workspace = true diff --git a/src/cli/src/data/export_v2.rs b/src/cli/src/data/export_v2.rs index f9a3c2be095..9887e525030 100644 --- a/src/cli/src/data/export_v2.rs +++ b/src/cli/src/data/export_v2.rs @@ -41,11 +41,13 @@ //! `--experimental-metric-export` enables shared Metric physical scans for Parquet. //! Enable `experimental_metric_export = true` on the frontend or standalone server //! and use an endpoint whose frontends all support and enable this option. +//! Add `--metric-data-layout packed` to write version-2 packed snapshots. +//! This requires the server packed-export capability; schema-only exports remain version 1. //! The command checks the server capability before modifying the snapshot. //! Each snapshot path belongs to one export task. Resume requires stable source //! schemas/data and confirmation that the previous export and storage writes ended; //! an HTTP timeout does not establish that. Only unfinished chunks are cleaned and -//! rerun. Completed chunks and the V2 snapshot/import format remain unchanged. +//! rerun. Resume retains the snapshot version/layout and completed chunks. mod chunker; mod command; diff --git a/src/cli/src/data/export_v2/command.rs b/src/cli/src/data/export_v2/command.rs index e0f67eb796c..20e69ef5a99 100644 --- a/src/cli/src/data/export_v2/command.rs +++ b/src/cli/src/data/export_v2/command.rs @@ -296,6 +296,10 @@ pub struct ExportCreateCommand { #[clap(long)] experimental_metric_export: bool, + /// Store independent Metric Parquet streams in shared packed objects. + #[clap(long, value_parser = ["packed"], requires = "experimental_metric_export")] + metric_data_layout: Option, + /// Delete existing snapshot and recreate. #[clap(long)] force: bool, @@ -345,6 +349,16 @@ impl ExportCreateCommand { let time_range = TimeRange::parse(self.start_time.as_deref(), self.end_time.as_deref()) .map_err(BoxedError::new)?; + if self.metric_data_layout.is_some() + && time_range.is_bounded() + && time_range.start == time_range.end + { + return crate::error::InvalidArgumentsSnafu { + msg: "Packed export requires --start-time to be earlier than --end-time", + } + .fail() + .map_err(BoxedError::new); + } if self.chunk_time_window.is_some() && !time_range.is_bounded() { return ChunkTimeWindowRequiresBoundsSnafu .fail() @@ -425,6 +439,13 @@ impl ExportCreateCommand { .map_err(BoxedError::new); } } + if self.metric_data_layout.is_some() && !self.schema_only { + database_client + .require_packed_export() + .await + .context(DatabaseSnafu) + .map_err(BoxedError::new)?; + } let storage = OpenDalStorage::from_uri(&self.to, &self.storage).map_err(BoxedError::new)?; Ok(Box::new(ExportCreate { @@ -434,6 +455,7 @@ impl ExportCreateCommand { schema_only: self.schema_only, format: self.format, experimental_metric_export: self.experimental_metric_export, + packed: self.metric_data_layout.is_some(), force: self.force, time_range, chunk_time_window: self.chunk_time_window, @@ -462,6 +484,7 @@ struct ExportConfig { schema_only: bool, format: DataFormat, experimental_metric_export: bool, + packed: bool, force: bool, time_range: TimeRange, chunk_time_window: Option, @@ -504,7 +527,7 @@ impl ExportCreate { let mut manifest = self.storage.read_manifest().await?; // Check version compatibility - if manifest.version != MANIFEST_VERSION || manifest.data_layout.is_some() { + if manifest.validate_layout().is_err() { return ManifestVersionMismatchSnafu { expected: MANIFEST_VERSION, found: manifest.version, @@ -513,6 +536,15 @@ impl ExportCreate { } validate_resume_config(&manifest, &self.config)?; + if manifest.is_packed() { + self.database_client + .require_packed_export() + .await + .context(DatabaseSnafu)?; + if !self.config.experimental_metric_export { + return crate::data::export_v2::error::MetricExportUnavailableSnafu.fail(); + } + } info!( "Resuming existing snapshot: {} (completed: {}/{} chunks)", @@ -571,10 +603,14 @@ impl ExportCreate { self.config.chunk_time_window, )?; + if self.config.packed && !self.config.schema_only { + manifest.version = 2; + manifest.data_layout = Some(common_datasource::packed_snapshot::PACKED_LAYOUT.into()); + } if self.config.experimental_metric_export { for chunk in &manifest.chunks { self.storage - .prepare_export_chunk(&schema_names, chunk.id, false) + .prepare_export_chunk(&schema_names, chunk.id, false, manifest.is_packed()) .await?; } } @@ -796,6 +832,14 @@ fn build_schema_ddl( } fn validate_resume_config(manifest: &Manifest, config: &ExportConfig) -> Result<()> { + if config.packed && !manifest.schema_only && !manifest.is_packed() { + return ResumeConfigMismatchSnafu { + field: "metric_data_layout", + existing: "standalone".to_string(), + requested: "packed".to_string(), + } + .fail(); + } if manifest.schema_only != config.schema_only { return SchemaOnlyModeMismatchSnafu { existing_schema_only: manifest.schema_only, @@ -1801,6 +1845,7 @@ mod tests { schema_only: false, format: DataFormat::Parquet, experimental_metric_export: false, + packed: false, force: false, time_range: TimeRange::unbounded(), chunk_time_window: None, @@ -1820,7 +1865,7 @@ mod tests { #[test] fn test_validate_resume_config_accepts_schema_selection_with_different_case_and_order() { - let manifest = Manifest::new_for_export( + let mut manifest = Manifest::new_for_export( "greptime".to_string(), vec!["public".to_string(), "analytics".to_string()], false, @@ -1829,7 +1874,7 @@ mod tests { None, ) .unwrap(); - let config = ExportConfig { + let mut config = ExportConfig { catalog: "greptime".to_string(), schemas: Some(vec![ "ANALYTICS".to_string(), @@ -1839,6 +1884,7 @@ mod tests { schema_only: false, format: DataFormat::Parquet, experimental_metric_export: false, + packed: false, force: false, time_range: TimeRange::unbounded(), chunk_time_window: None, @@ -1850,6 +1896,18 @@ mod tests { }; assert!(validate_resume_config(&manifest, &config).is_ok()); + config.packed = true; + assert!( + validate_resume_config(&manifest, &config) + .unwrap_err() + .to_string() + .contains("metric_data_layout") + ); + manifest.version = 2; + manifest.data_layout = Some(common_datasource::packed_snapshot::PACKED_LAYOUT.into()); + assert!(validate_resume_config(&manifest, &config).is_ok()); + config.packed = false; + assert!(validate_resume_config(&manifest, &config).is_ok()); } #[test] @@ -1872,6 +1930,7 @@ mod tests { schema_only: false, format: DataFormat::Parquet, experimental_metric_export: false, + packed: false, force: false, time_range, chunk_time_window: Some(Duration::from_secs(3600)), @@ -1906,6 +1965,7 @@ mod tests { schema_only: false, format: DataFormat::Csv, experimental_metric_export: false, + packed: false, force: false, time_range: TimeRange::unbounded(), chunk_time_window: None, @@ -1942,6 +2002,7 @@ mod tests { schema_only: false, format: DataFormat::Parquet, experimental_metric_export: false, + packed: false, force: false, time_range: TimeRange::new(Some(start), Some(start)), chunk_time_window: None, diff --git a/src/cli/src/data/export_v2/coordinator.rs b/src/cli/src/data/export_v2/coordinator.rs index 5f95c43603e..c6181e51211 100644 --- a/src/cli/src/data/export_v2/coordinator.rs +++ b/src/cli/src/data/export_v2/coordinator.rs @@ -40,6 +40,7 @@ struct ExportContext<'a> { format: DataFormat, parallelism: usize, experimental_metric_export: bool, + packed: bool, resume: bool, } @@ -71,6 +72,7 @@ pub async fn export_data( catalog: manifest.catalog.clone(), schemas: manifest.schemas.clone(), format: manifest.format, + packed: manifest.is_packed(), parallelism: options.parallelism, experimental_metric_export: options.experimental_metric_export, resume: options.resume, @@ -279,7 +281,7 @@ async fn export_chunk( if context.experimental_metric_export { context .storage - .prepare_export_chunk(&context.schemas, chunk_id, context.resume) + .prepare_export_chunk(&context.schemas, chunk_id, context.resume, context.packed) .await?; } let scheme = StorageScheme::from_uri(context.snapshot_uri)?; @@ -289,6 +291,7 @@ async fn export_chunk( time_range, parallelism: context.parallelism, experimental_metric_export: context.experimental_metric_export, + packed: context.packed, }; for schema in &context.schemas { @@ -314,6 +317,39 @@ async fn export_chunk( } let files = list_chunk_files(context.storage, &context.schemas, chunk_id).await?; + if context.packed { + use common_datasource::packed_snapshot::{PACK_INDEX_FILE, PackIndex}; + let mut expected = Vec::new(); + for schema in &context.schemas { + let prefix = data_dir_for_schema_chunk(schema, chunk_id); + let index_path = format!("{prefix}{PACK_INDEX_FILE}"); + let text = context.storage.read_text(&index_path).await?; + let index: PackIndex = serde_json::from_str(&text).map_err(|e| { + crate::data::export_v2::error::InvalidUriSnafu { + uri: context.snapshot_uri, + reason: e.to_string(), + } + .build() + })?; + index.validate().map_err(|e| { + crate::data::export_v2::error::InvalidUriSnafu { + uri: context.snapshot_uri, + reason: e.to_string(), + } + .build() + })?; + expected.push(index_path); + expected.extend(index.objects.iter().map(|o| format!("{prefix}{}", o.path))); + } + expected.sort(); + if files != expected { + return crate::data::export_v2::error::InvalidUriSnafu { + uri: context.snapshot_uri, + reason: "chunk inventory differs from packed indexes", + } + .fail(); + } + } info!("Collected {} files for chunk {}", files.len(), chunk_id); Ok(files) } @@ -462,6 +498,158 @@ mod tests { } } + struct PausedClose { + inner: Option, + started: std::sync::Arc, + release: std::sync::Arc, + } + + impl object_store::layers::mock::oio::Write for PausedClose { + async fn write(&mut self, bytes: object_store::Buffer) -> object_store::Result<()> { + self.inner.as_mut().unwrap().write(bytes).await + } + + async fn close(&mut self) -> object_store::Result { + let mut inner = self.inner.take().unwrap(); + let (started, release) = (self.started.clone(), self.release.clone()); + // Model blocking storage I/O that survives dropping its caller. + let (inner, result) = tokio::spawn(async move { + started.notify_one(); + release.notified().await; + let result = inner.close().await; + (inner, result) + }) + .await + .unwrap(); + self.inner = Some(inner); + result + } + + async fn abort(&mut self) -> object_store::Result<()> { + self.inner.as_mut().unwrap().abort().await + } + } + + #[tokio::test] + async fn packed_cancellation_drains_close_before_failing_chunk() { + use std::sync::Arc; + + use common_datasource::packed_snapshot::PACK_INDEX_FILE; + use common_datasource::packed_writer::{PackedTableWriter, PackedWriter}; + use common_datasource::parquet_writer::ParquetFileWriter; + use datatypes::arrow::datatypes::{DataType, Field, Schema}; + use object_store::layers::mock::{MockLayerBuilder, MockWriterFactory}; + use tokio_util::sync::CancellationToken; + + use crate::data::snapshot_storage::OpenDalStorage; + + for paused_path in ["pack-000000.bin", PACK_INDEX_FILE] { + let directory = tempfile::tempdir().unwrap(); + let uri = url::Url::from_directory_path(directory.path()).unwrap(); + let storage = + OpenDalStorage::from_uri(uri.as_str(), &ObjectStoreConfig::default()).unwrap(); + let started = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let factory: MockWriterFactory = Arc::new({ + let (started, release) = (started.clone(), release.clone()); + move |path, _, inner| { + if path.trim_start_matches('/') == paused_path { + Box::new(PausedClose { + inner: Some(inner), + started: started.clone(), + release: release.clone(), + }) + } else { + inner + } + } + }); + let store = object_store::secure_fs::SecureFsRoot::open(directory.path()) + .unwrap() + .build_operator() + .layer( + MockLayerBuilder::default() + .writer_factory(factory) + .build() + .unwrap(), + ); + let packed = PackedWriter::new(store.clone()).unwrap(); + let schema = Arc::new(Schema::new(vec![Field::new( + "value", + DataType::Int64, + true, + )])); + let mut table = ParquetFileWriter::open_packed( + schema, + store.clone(), + "unused", + None, + PackedTableWriter::new(packed.clone(), "empty".into(), 1, false), + ) + .unwrap(); + table.finish(None).await.unwrap(); + let token = CancellationToken::new(); + let mut manifest = pending_manifest(1); + manifest.version = 2; + manifest.data_layout = Some(common_datasource::packed_snapshot::PACKED_LAYOUT.into()); + let export = export_data_concurrent( + &storage, + &mut manifest, + 2, + &crate::data::progress::NoopProgress, + |_, _| async { + let mut packed = packed.lock().await; + let result = packed.finish(&token).await; + if result.is_err() { + packed.abort().await.unwrap(); + } + result.map_err(|error| { + crate::data::export_v2::error::IoSnafu { + operation: "exporting packed chunk", + error: std::io::Error::other(error), + } + .build() + }) + }, + ); + tokio::pin!(export); + tokio::select! { + _ = started.notified() => {}, + result = &mut export => panic!("completed while close paused: {result:?}"), + } + token.cancel(); + assert!(futures::poll!(&mut export).is_pending()); + assert_eq!( + storage.read_manifest().await.unwrap().chunks[0].status, + ChunkStatus::InProgress + ); + release.notify_one(); + assert!(export.await.unwrap_err().to_string().contains("cancelled")); + let persisted = storage.read_manifest().await.unwrap(); + assert_eq!(persisted.chunks[0].status, ChunkStatus::Failed); + assert!(!persisted.is_complete()); + assert!(persisted.chunks[0].files.is_empty()); + assert!( + store + .stat("pack-000000.bin") + .await + .unwrap() + .content_length() + > 0 + ); + assert_eq!( + store.exists(PACK_INDEX_FILE).await.unwrap(), + paused_path == PACK_INDEX_FILE + ); + if paused_path == PACK_INDEX_FILE { + let index: common_datasource::packed_snapshot::PackIndex = + serde_json::from_slice(&store.read(PACK_INDEX_FILE).await.unwrap().to_bytes()) + .unwrap(); + index.validate_membership(["empty"]).unwrap(); + } + } + } + #[test] fn test_next_eligible_chunk_scans_in_order() { let manifest = pending_manifest(3); diff --git a/src/cli/src/data/export_v2/data.rs b/src/cli/src/data/export_v2/data.rs index 8723154e4b1..53c3899f59e 100644 --- a/src/cli/src/data/export_v2/data.rs +++ b/src/cli/src/data/export_v2/data.rs @@ -31,6 +31,7 @@ pub(super) struct CopyOptions { pub(super) time_range: TimeRange, pub(super) parallelism: usize, pub(crate) experimental_metric_export: bool, + pub(super) packed: bool, } pub(super) struct CopyTarget { @@ -220,6 +221,9 @@ fn build_with_options(options: &CopyOptions) -> String { if options.experimental_metric_export { parts.push("experimental_metric_export='true'".to_string()); } + if options.packed { + parts.push("metric_data_layout='packed'".to_string()); + } if let Some(start) = options.time_range.start { parts.push(format!( "START_TIME='{}'", diff --git a/src/cli/src/data/export_v2/manifest.rs b/src/cli/src/data/export_v2/manifest.rs index 5dd9d6f1684..d56f2732490 100644 --- a/src/cli/src/data/export_v2/manifest.rs +++ b/src/cli/src/data/export_v2/manifest.rs @@ -28,7 +28,7 @@ use crate::data::export_v2::error::{ TimeParseInvalidFormatSnafu, }; -/// Manifest format version produced by the current exporter. +/// Default standalone manifest version. Packed data snapshots use version 2. pub const MANIFEST_VERSION: u32 = 1; /// Manifest file name within snapshot directory. diff --git a/src/cli/src/data/snapshot_storage.rs b/src/cli/src/data/snapshot_storage.rs index c5e9561d285..17008706066 100644 --- a/src/cli/src/data/snapshot_storage.rs +++ b/src/cli/src/data/snapshot_storage.rs @@ -295,8 +295,9 @@ pub trait SnapshotStorage: Send + Sync { schemas: &[String], chunk_id: u32, resume: bool, + packed: bool, ) -> Result<()> { - let _ = (schemas, chunk_id, resume); + let _ = (schemas, chunk_id, resume, packed); InvalidUriSnafu { uri: "snapshot", reason: "storage does not support preparing export chunks", @@ -730,6 +731,7 @@ impl SnapshotStorage for OpenDalStorage { schemas: &[String], chunk_id: u32, resume: bool, + packed: bool, ) -> Result<()> { let mut files = Vec::new(); for schema in schemas { @@ -751,10 +753,17 @@ impl SnapshotStorage for OpenDalStorage { continue; } let name = path.strip_prefix(&prefix).unwrap_or(""); - if !resume || entry.metadata().is_dir() || !valid_chunk_filename(name) { + if !resume + || entry.metadata().is_dir() + || !(if packed { + valid_packed_chunk_filename(name) + } else { + valid_chunk_filename(name) + }) + { return InvalidUriSnafu { uri: path, - reason: "expected an empty new chunk or direct Parquet files in an owned unfinished chunk", + reason: "expected an empty new chunk or recognized files in an owned unfinished chunk", }.fail(); } files.push(path.to_string()); @@ -786,6 +795,17 @@ impl SnapshotStorage for OpenDalStorage { } } +fn valid_packed_chunk_filename(name: &str) -> bool { + name == common_datasource::packed_snapshot::PACK_INDEX_FILE + || [("pack-", ".bin"), ("table-", ".parquet")] + .iter() + .any(|(prefix, suffix)| { + name.strip_prefix(prefix) + .and_then(|n| n.strip_suffix(suffix)) + .is_some_and(|id| !id.is_empty() && id.bytes().all(|b| b.is_ascii_digit())) + }) +} + fn valid_chunk_filename(name: &str) -> bool { name.strip_suffix(".parquet") .is_some_and(|stem| !stem.is_empty()) @@ -1074,7 +1094,7 @@ mod tests { let schemas = vec!["public".to_string(), "other".to_string()]; assert!( storage - .prepare_export_chunk(&schemas, 2, false) + .prepare_export_chunk(&schemas, 2, false, false) .await .is_err() ); @@ -1085,11 +1105,11 @@ mod tests { .unwrap() ); storage - .prepare_export_chunk(&schemas, 2, true) + .prepare_export_chunk(&schemas, 2, true, false) .await .unwrap(); storage - .prepare_export_chunk(&schemas, 2, true) + .prepare_export_chunk(&schemas, 2, true, false) .await .unwrap(); assert!( @@ -1129,7 +1149,7 @@ mod tests { } assert!( storage - .prepare_export_chunk(&["public".into(), "other".into()], 2, true) + .prepare_export_chunk(&["public".into(), "other".into()], 2, true, false) .await .is_err() ); @@ -1143,6 +1163,70 @@ mod tests { } } + #[tokio::test] + async fn packed_retry_preserves_foreign_files_and_completed_chunks() { + let dir = tempdir().unwrap(); + let storage = make_storage_with_rooted_fs(dir.path()); + let owned = ["pack-000000.bin", "table-42.parquet", "pack-index.json"]; + for name in owned { + storage + .write_text(&format!("data/public/2/{name}"), "partial") + .await + .unwrap(); + } + for path in [ + "data/public/1/pack-000000.bin", + "data/public/2/foreign.parquet", + ] { + storage.write_text(path, "keep").await.unwrap(); + } + let schemas = ["public".into()]; + assert!( + storage + .prepare_export_chunk(&schemas, 2, true, true) + .await + .is_err() + ); + for name in owned { + assert!( + storage + .file_exists(&format!("data/public/2/{name}")) + .await + .unwrap() + ); + } + assert_eq!( + storage + .read_text("data/public/2/foreign.parquet") + .await + .unwrap(), + "keep" + ); + storage + .object_store + .delete("data/public/2/foreign.parquet") + .await + .unwrap(); + storage + .prepare_export_chunk(&schemas, 2, true, true) + .await + .unwrap(); + assert!( + storage + .list_files_recursive("data/public/2/") + .await + .unwrap() + .is_empty() + ); + assert_eq!( + storage + .read_text("data/public/1/pack-000000.bin") + .await + .unwrap(), + "keep" + ); + } + #[tokio::test] async fn test_prepare_export_chunk_reports_delete_failure() { let dir = tempdir().unwrap(); @@ -1159,7 +1243,7 @@ mod tests { .await .unwrap(); let error = storage - .prepare_export_chunk(&["public".into()], 2, true) + .prepare_export_chunk(&["public".into()], 2, true, false) .await .unwrap_err(); let Error::StorageOperation { @@ -1191,7 +1275,7 @@ mod tests { .await .unwrap(); storage - .prepare_export_chunk(&[schema], 2, true) + .prepare_export_chunk(&[schema], 2, true, false) .await .unwrap(); assert!(!storage.file_exists(&owned).await.unwrap()); diff --git a/src/cli/src/database.rs b/src/cli/src/database.rs index b771ab48d05..ac314bc5899 100644 --- a/src/cli/src/database.rs +++ b/src/cli/src/database.rs @@ -99,6 +99,15 @@ impl DatabaseClient { /// Requires the explicit packed-import protocol before any restore mutation. pub async fn require_packed_import(&self) -> Result<()> { + self.require_packed_capability("metric_packed_import").await + } + + /// Requires packed export support before creating snapshot artifacts. + pub async fn require_packed_export(&self) -> Result<()> { + self.require_packed_capability("metric_packed_export").await + } + + async fn require_packed_capability(&self, capability: &str) -> Result<()> { let url = format!("http://{}/v1/capabilities", self.addr); let mut builder = reqwest::Client::builder().timeout(self.timeout); if let Some(proxy) = self.proxy.clone() { @@ -113,18 +122,18 @@ impl DatabaseClient { request = request.header("Authorization", auth); } let response = request.send().await.with_context(|_| HttpQuerySqlSnafu { - reason: "packed import capability request failed", + reason: "packed snapshot capability request failed", })?; if response.status() == reqwest::StatusCode::NOT_FOUND { return crate::error::InvalidArgumentsSnafu { - msg: "target does not support packed import", + msg: format!("server does not support {capability}"), } .fail(); } let response = response .error_for_status() .with_context(|_| HttpQuerySqlSnafu { - reason: "packed import capability request rejected", + reason: "packed snapshot capability request rejected", })?; let body = response.text().await.with_context(|_| HttpQuerySqlSnafu { reason: "cannot read capability response", @@ -136,9 +145,9 @@ impl DatabaseClient { } .fail(); } - if value.get("metric_packed_import").and_then(Value::as_u64) != Some(1) { + if value.get(capability).and_then(Value::as_u64) != Some(1) { return crate::error::InvalidArgumentsSnafu { - msg: "target does not support packed import version 1", + msg: format!("server does not support {capability} version 1"), } .fail(); } @@ -221,41 +230,54 @@ mod tests { #[tokio::test] async fn packed_capability_probe_checks_auth_and_protocol_version() { use tokio::io::{AsyncReadExt, AsyncWriteExt}; - for (status, body, supported) in [ - (200, r#"{"metric_packed_import":1,"future":2}"#, true), - (200, "{}", false), - (200, r#"{"metric_packed_import":2}"#, false), - (404, "{}", false), - (401, "{}", false), - (403, "{}", false), - (200, "[]", false), - (200, "invalid-json", false), - ] { - let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - let address = listener.local_addr().unwrap(); - let server = tokio::spawn(async move { - let (mut socket, _) = listener.accept().await.unwrap(); - let mut request = vec![0; 4096]; - let n = socket.read(&mut request).await.unwrap(); - let request = String::from_utf8_lossy(&request[..n]).to_lowercase(); - assert!(request.starts_with("get /v1/capabilities ")); - assert!(request.contains("authorization: basic dxnlcjpwyxnzd29yza==")); - socket.write_all(format!("HTTP/1.1 {status} Response\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len()).as_bytes()).await.unwrap(); - }); - let client = super::DatabaseClient::new( - address.to_string(), - "greptime".into(), - Some("user:password".into()), - std::time::Duration::from_secs(5), - None, - true, - ); - assert_eq!( - client.require_packed_import().await.is_ok(), - supported, - "{status}: {body}" - ); - server.await.unwrap(); + for export in [false, true] { + for (status, body, supported) in [ + (200, r#"{"metric_packed_import":1,"future":2}"#, true), + (200, "{}", false), + (200, r#"{"metric_packed_import":2}"#, false), + (404, "{}", false), + (401, "{}", false), + (403, "{}", false), + (200, "[]", false), + (200, "invalid-json", false), + ] { + let body = if export { + body.replace("metric_packed_import", "metric_packed_export") + } else { + body.to_string() + }; + let response_body = body.clone(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let body = response_body; + let (mut socket, _) = listener.accept().await.unwrap(); + let mut request = vec![0; 4096]; + let n = socket.read(&mut request).await.unwrap(); + let request = String::from_utf8_lossy(&request[..n]).to_lowercase(); + assert!(request.starts_with("get /v1/capabilities ")); + assert!(request.contains("authorization: basic dxnlcjpwyxnzd29yza==")); + socket.write_all(format!("HTTP/1.1 {status} Response\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", body.len()).as_bytes()).await.unwrap(); + }); + let client = super::DatabaseClient::new( + address.to_string(), + "greptime".into(), + Some("user:password".into()), + std::time::Duration::from_secs(5), + None, + true, + ); + assert_eq!( + if export { + client.require_packed_export().await.is_ok() + } else { + client.require_packed_import().await.is_ok() + }, + supported, + "{status}: {body}" + ); + server.await.unwrap(); + } } } diff --git a/src/common/datasource/Cargo.toml b/src/common/datasource/Cargo.toml index d1ca609c606..3953d144bad 100644 --- a/src/common/datasource/Cargo.toml +++ b/src/common/datasource/Cargo.toml @@ -39,6 +39,7 @@ parquet.workspace = true paste.workspace = true regex.workspace = true serde.workspace = true +serde_json.workspace = true snafu.workspace = true strum.workspace = true tokio.workspace = true diff --git a/src/common/datasource/src/lib.rs b/src/common/datasource/src/lib.rs index c832831da80..ebec28d7129 100644 --- a/src/common/datasource/src/lib.rs +++ b/src/common/datasource/src/lib.rs @@ -20,6 +20,7 @@ pub mod file_format; pub mod lister; pub mod object_store; pub mod packed_snapshot; +pub mod packed_writer; pub mod parquet_writer; pub mod share_buffer; #[cfg(test)] diff --git a/src/common/datasource/src/packed_writer.rs b/src/common/datasource/src/packed_writer.rs new file mode 100644 index 00000000000..635a1a7a771 --- /dev/null +++ b/src/common/datasource/src/packed_writer.rs @@ -0,0 +1,669 @@ +// Copyright 2023 Greptime Team +// +// Licensed 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. + +//! Streaming destinations for one schema chunk of independent Parquet streams. + +use std::sync::Arc; + +use bytes::Bytes; +use object_store::{ObjectStore, Writer}; +use snafu::{ResultExt, ensure}; +use tokio::sync::Mutex; +use tokio_util::sync::CancellationToken; + +use crate::error::{self, Result}; +use crate::packed_snapshot::{ObjectKind, PACK_INDEX_FILE, PackIndex, PackObject, PackTable}; + +pub const WRITE_BYTES: usize = 8 * 1024 * 1024; + +/// One request owns this sink. Finishers retain their encoder admission while +/// waiting on the mutex, bounding both completed buffers and pending appends. +pub type PackedWriterRef = Arc>; + +#[derive(Clone, Copy)] +struct UploadLimits { + part: usize, + object: u64, +} + +impl UploadLimits { + fn for_store(store: &ObjectStore) -> Result { + let info = store.info(); + let caps = info.capability(); + let part = caps + .write_multi_max_size + .unwrap_or(WRITE_BYTES) + .min(WRITE_BYTES); + ensure!( + part > 0 && caps.write_multi_min_size.unwrap_or(0) <= part, + error::ParquetWriterResourceSnafu { + reason: "backend cannot accept bounded 8 MiB parts" + } + ); + // OpenDAL exposes part size but not part count. These are the limits of + // its multipart implementations, including GCS's XML multipart writer. + let parts = match info.scheme() { + "s3" | "oss" | "gcs" => 10_000, + "azblob" => 50_000, + _ => u64::MAX, + }; + let object = (caps.write_total_max_size.unwrap_or(usize::MAX) as u64) + .min(parts.saturating_mul(part as u64)); + Ok(Self { part, object }) + } +} + +struct Upload { + writer: Writer, + store: ObjectStore, + path: String, + limits: UploadLimits, + length: u64, + conditional: bool, + close_started: bool, +} + +impl Upload { + async fn open(store: &ObjectStore, path: String, limits: UploadLimits) -> Result { + let conditional = store.info().capability().write_with_if_not_exists; + if !conditional { + ensure!( + !store + .exists(&path) + .await + .context(error::ReadObjectSnafu { path: &path })?, + error::InvalidPackedSnapshotSnafu { + reason: format!("output already exists: {path}") + } + ); + } + let writer = store + .writer_with(&path) + .chunk(limits.part) + .concurrent(1) + .if_not_exists(conditional) + .await + .context(error::WriteObjectSnafu { path: &path })?; + Ok(Self { + writer, + store: store.clone(), + path, + limits, + length: 0, + conditional, + close_started: false, + }) + } + + async fn write(&mut self, bytes: Bytes) -> Result<()> { + ensure!( + (bytes.len() as u64) <= self.limits.object.saturating_sub(self.length), + error::ParquetWriterResourceSnafu { + reason: "packed export object or multipart part-count limit exceeded" + } + ); + for offset in (0..bytes.len()).step_by(self.limits.part) { + self.writer + .write(bytes.slice(offset..(offset + self.limits.part).min(bytes.len()))) + .await + .context(error::WriteObjectSnafu { path: &self.path })?; + } + self.length += bytes.len() as u64; + Ok(()) + } + + async fn write_json_array( + &mut self, + values: &[T], + token: &CancellationToken, + ) -> Result<()> { + self.write(Bytes::from_static(b"[")).await?; + for (i, value) in values.iter().enumerate() { + check_cancelled(Some(token))?; + if i > 0 { + self.write(Bytes::from_static(b",")).await?; + } + let bytes = serde_json::to_vec(value).map_err(|e| { + error::InvalidPackedSnapshotSnafu { + reason: e.to_string(), + } + .build() + })?; + self.write(bytes.into()).await?; + } + self.write(Bytes::from_static(b"]")).await + } + + async fn close(&mut self) -> Result<()> { + self.close_started = true; + self.writer + .close() + .await + .context(error::WriteObjectSnafu { path: &self.path })?; + Ok(()) + } + + async fn abort(mut self) -> Result<()> { + let result = self.writer.abort().await; + if !self.conditional + && result.as_ref().is_err_and(|e| { + e.kind() == object_store::ErrorKind::Unsupported + && (!self.close_started + || object_store::secure_fs::is_unsynced_overwrite_abort(e)) + }) + { + let store = self.store.clone(); + let path = self.path.clone(); + drop(self); + store + .delete(&path) + .await + .context(error::WriteObjectSnafu { path })?; + } else { + result.context(error::WriteObjectSnafu { path: &self.path })?; + } + Ok(()) + } +} + +/// Pack/index metadata is accounted separately from retained Arrow payloads. +pub struct PackedWriter { + store: ObjectStore, + limits: UploadLimits, + pack: Option, + index: PackIndex, + next_pack: usize, +} + +impl PackedWriter { + pub fn new(store: ObjectStore) -> Result { + let limits = UploadLimits::for_store(&store)?; + Ok(Arc::new(Mutex::new(Self { + store, + limits, + pack: None, + index: PackIndex { + version: 1, + objects: vec![], + tables: vec![], + }, + next_pack: 0, + }))) + } + + async fn close_pack(&mut self) -> Result<()> { + if let Some(pack) = &mut self.pack { + pack.close().await?; + self.index.objects.push(PackObject { + path: pack.path.clone(), + kind: ObjectKind::Pack, + length: pack.length, + }); + self.pack = None; + } + Ok(()) + } + + async fn append( + &mut self, + name: String, + bytes: Bytes, + rows: u64, + token: Option<&CancellationToken>, + ) -> Result<()> { + check_cancelled(token)?; + if self + .pack + .as_ref() + .is_some_and(|p| bytes.len() as u64 > self.limits.object.saturating_sub(p.length)) + { + self.close_pack().await?; + } + check_cancelled(token)?; + if self.pack.is_none() { + let path = format!("pack-{:06}.bin", self.next_pack); + self.next_pack += 1; + self.pack = Some(Upload::open(&self.store, path, self.limits).await?); + } + let pack = self.pack.as_mut().ok_or_else(|| { + error::InvalidPackedSnapshotSnafu { + reason: "missing pack writer", + } + .build() + })?; + let table = PackTable { + table_name: name, + object: pack.path.clone(), + offset: pack.length, + length: bytes.len() as u64, + row_count: rows, + }; + pack.write(bytes).await?; + check_cancelled(token)?; + self.index.tables.push(table); + Ok(()) + } + + /// Publish the index only after every data object has closed successfully. + pub async fn finish(&mut self, token: &CancellationToken) -> Result> { + check_cancelled(Some(token))?; + self.close_pack().await?; + check_cancelled(Some(token))?; + self.index + .tables + .sort_unstable_by(|a, b| a.table_name.cmp(&b.table_name)); + self.index.validate()?; + let mut upload = Upload::open(&self.store, PACK_INDEX_FILE.into(), self.limits).await?; + let result = async { + upload + .write(Bytes::from_static(b"{\"version\":1,\"objects\":")) + .await?; + upload.write_json_array(&self.index.objects, token).await?; + upload.write(Bytes::from_static(b",\"tables\":")).await?; + upload.write_json_array(&self.index.tables, token).await?; + upload.write(Bytes::from_static(b"}")).await?; + check_cancelled(Some(token))?; + upload.close().await?; + check_cancelled(Some(token)) + } + .await; + if let Err(error) = result { + if let Err(secondary) = upload.abort().await { + common_telemetry::warn!(secondary; "Failed to abort pack index"); + } + return Err(error); + } + Ok(self + .index + .objects + .iter() + .map(|o| o.path.clone()) + .chain([PACK_INDEX_FILE.into()]) + .collect()) + } + + /// Call only after all admitted encoders have drained their I/O. + pub async fn abort(&mut self) -> Result<()> { + if let Some(pack) = self.pack.take() { + pack.abort().await?; + } + Ok(()) + } +} + +/// Buffers a small stream, or spills its prefix and continues the same stream. +pub struct PackedTableWriter { + shared: PackedWriterRef, + name: String, + path: String, + buffer: Vec, + ordinary: bool, + standalone: Option, +} + +impl PackedTableWriter { + /// The caller retains encoder admission until finish/abort returns. + pub fn new(shared: PackedWriterRef, name: String, id: u32, ordinary: bool) -> Self { + Self { + shared, + name, + path: format!("table-{id}.parquet"), + buffer: if ordinary { + Vec::new() + } else { + Vec::with_capacity(WRITE_BYTES) + }, + ordinary, + standalone: None, + } + } + + pub(crate) async fn write(&mut self, bytes: Bytes) -> Result<()> { + if self.standalone.is_none() + && (self.ordinary || self.buffer.len().saturating_add(bytes.len()) > WRITE_BYTES) + { + let sink = self.shared.lock().await; + self.standalone = + Some(Upload::open(&sink.store, self.path.clone(), sink.limits).await?); + } + if let Some(upload) = &mut self.standalone { + if !self.buffer.is_empty() { + upload + .write(std::mem::take(&mut self.buffer).into()) + .await?; + } + upload.write(bytes).await + } else { + self.buffer.extend_from_slice(&bytes); + Ok(()) + } + } + + pub(crate) async fn finish( + &mut self, + rows: u64, + token: Option<&CancellationToken>, + ) -> Result<()> { + check_cancelled(token)?; + if let Some(upload) = &mut self.standalone { + upload.close().await?; + check_cancelled(token)?; + let mut shared = self.shared.lock().await; + shared.index.objects.push(PackObject { + path: self.path.clone(), + kind: ObjectKind::Parquet, + length: upload.length, + }); + shared.index.tables.push(PackTable { + table_name: self.name.clone(), + object: self.path.clone(), + offset: 0, + length: upload.length, + row_count: rows, + }); + } else { + self.shared + .lock() + .await + .append( + self.name.clone(), + // OpenDAL may retain each slice until the pack part fills. + // Release the encoder's 8 MiB capacity before queuing it. + Bytes::from(std::mem::take(&mut self.buffer).into_boxed_slice()), + rows, + token, + ) + .await?; + } + Ok(()) + } + + pub(crate) async fn abort(&mut self) -> Result<()> { + if let Some(upload) = self.standalone.take() { + upload.abort().await?; + } + Ok(()) + } +} + +fn check_cancelled(token: Option<&CancellationToken>) -> Result<()> { + ensure!( + token.is_none_or(|t| !t.is_cancelled()), + error::ParquetWriteCancelledSnafu + ); + Ok(()) +} + +#[cfg(test)] +mod tests { + use arrow::array::{ArrayRef, Int64Array, StringArray}; + use arrow::record_batch::RecordBatch; + use futures::TryStreamExt; + use object_store::layers::mock::{Metadata, MockLayerBuilder, MockWriterFactory, oio}; + use parquet::arrow::ParquetRecordBatchStreamBuilder; + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; + + use super::*; + use crate::file_format::parquet::packed_reader::{PackReadWindows, PackedParquetReader}; + use crate::parquet_writer::{ParquetFileWriter, ParquetWriterLimits}; + + #[tokio::test] + async fn independent_streams_empty_and_spilled_roundtrip() { + let directory = common_test_util::temp_dir::create_temp_dir("packed-writer"); + let store = object_store::secure_fs::SecureFsRoot::open(directory.path()) + .unwrap() + .build_operator(); + let sizes = Arc::new(std::sync::Mutex::new(Vec::new())); + let factory: MockWriterFactory = Arc::new({ + let sizes = sizes.clone(); + move |_, _, writer| Box::new(Observe(writer, sizes.clone())) + }); + let store = store.layer( + MockLayerBuilder::default() + .writer_factory(factory) + .build() + .unwrap(), + ); + let shared = PackedWriter::new(store.clone()).unwrap(); + let small = RecordBatch::try_from_iter([( + "tag", + Arc::new(StringArray::from(vec![Some(""), None, Some("value")])) as ArrayRef, + )]) + .unwrap(); + let mut state = 17u64; + let values = (0..600_000) + .map(|_| { + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + state as i64 + }) + .collect::>(); + let random = Arc::new(Int64Array::from(values)) as ArrayRef; + let large = RecordBatch::try_from_iter([("a", random.clone()), ("b", random)]).unwrap(); + let batches = [small.clone(), RecordBatch::new_empty(large.schema()), large]; + for (id, batch) in batches.iter().enumerate() { + let sink = + PackedTableWriter::new(shared.clone(), format!("table{id}"), id as u32, false); + let mut writer = ParquetFileWriter::open_packed( + batch.schema(), + store.clone(), + "unused", + Some(ParquetWriterLimits { + row_group_rows: 8192, + flush_threshold_bytes: WRITE_BYTES, + max_row_groups: 4096, + }), + sink, + ) + .unwrap(); + writer.write(batch.clone(), None).await.unwrap(); + writer.finish(None).await.unwrap(); + } + let inventory = shared + .lock() + .await + .finish(&CancellationToken::new()) + .await + .unwrap(); + let index: PackIndex = + serde_json::from_slice(&store.read(PACK_INDEX_FILE).await.unwrap().to_bytes()).unwrap(); + index + .validate_membership(["table0", "table1", "table2"]) + .unwrap(); + assert_eq!(inventory.len(), 3); + assert!(sizes.lock().unwrap().iter().all(|n| *n <= WRITE_BYTES)); + assert_eq!( + index + .objects + .iter() + .filter(|o| o.kind == ObjectKind::Pack) + .count(), + 1 + ); + assert!( + index + .objects + .iter() + .any(|o| o.kind == ObjectKind::Parquet && o.length > WRITE_BYTES as u64) + ); + let windows = PackReadWindows::new(store.clone()); + for (entry, expected) in index.tables.iter().zip(batches) { + let object = index + .objects + .iter() + .find(|o| o.path == entry.object) + .unwrap(); + let actual = if object.kind == ObjectKind::Pack { + let reader = PackedParquetReader::new( + windows.clone(), + object.path.clone(), + object.length, + entry.offset, + entry.length, + ) + .unwrap(); + ParquetRecordBatchStreamBuilder::new(reader) + .await + .unwrap() + .build() + .unwrap() + .try_collect::>() + .await + .unwrap() + } else { + ParquetRecordBatchReaderBuilder::try_new( + store.read(&object.path).await.unwrap().to_bytes(), + ) + .unwrap() + .build() + .unwrap() + .collect::, _>>() + .unwrap() + }; + assert_eq!(entry.row_count, expected.num_rows() as u64); + assert_eq!( + arrow::compute::concat_batches(&expected.schema(), &actual).unwrap(), + expected + ); + } + } + + struct Observe(oio::Writer, Arc>>); + impl oio::Write for Observe { + async fn write(&mut self, bytes: object_store::Buffer) -> object_store::Result<()> { + self.1.lock().unwrap().push(bytes.len()); + self.0.write(bytes).await + } + async fn close(&mut self) -> object_store::Result { + self.0.close().await + } + async fn abort(&mut self) -> object_store::Result<()> { + self.0.abort().await + } + } + + #[tokio::test] + async fn rolls_at_table_boundary_and_aborts_standalone_at_part_limit() { + let directory = common_test_util::temp_dir::create_temp_dir("packed-limit"); + let sizes = Arc::new(std::sync::Mutex::new(Vec::new())); + let factory: MockWriterFactory = Arc::new({ + let sizes = sizes.clone(); + move |_, _, writer| Box::new(Observe(writer, sizes.clone())) + }); + let store = object_store::secure_fs::SecureFsRoot::open(directory.path()) + .unwrap() + .build_operator() + .layer( + MockLayerBuilder::default() + .writer_factory(factory) + .build() + .unwrap(), + ); + let shared = PackedWriter::new(store.clone()).unwrap(); + // Two data parts per object; no oversized fixture is needed. + shared.lock().await.limits = UploadLimits { + part: 512, + object: 1024, + }; + for id in 0..3 { + let mut table = PackedTableWriter::new(shared.clone(), format!("t{id}"), id, false); + table.write(vec![id as u8; 600].into()).await.unwrap(); + table.finish(0, None).await.unwrap(); + } + let mut table = PackedTableWriter::new(shared.clone(), "large".into(), 4, true); + table.write(vec![0; 1024].into()).await.unwrap(); + assert!(matches!( + table.write(Bytes::from_static(b"x")).await, + Err(error::Error::ParquetWriterResource { .. }) + )); + table.abort().await.unwrap(); + assert!(!store.exists("table-4.parquet").await.unwrap()); + shared + .lock() + .await + .finish(&CancellationToken::new()) + .await + .unwrap(); + let index: PackIndex = + serde_json::from_slice(&store.read(PACK_INDEX_FILE).await.unwrap().to_bytes()).unwrap(); + assert_eq!(index.objects.len(), 3); + assert!( + index + .tables + .iter() + .all(|t| t.offset == 0 && t.length == 600) + ); + assert!(sizes.lock().unwrap().iter().all(|n| *n <= 512)); + } + + struct FailClose(oio::Writer); + impl oio::Write for FailClose { + async fn write(&mut self, bytes: object_store::Buffer) -> object_store::Result<()> { + self.0.write(bytes).await + } + async fn close(&mut self) -> object_store::Result { + Err(object_store::Error::new( + object_store::ErrorKind::Unexpected, + "injected close failure", + )) + } + async fn abort(&mut self) -> object_store::Result<()> { + self.0.abort().await + } + } + + #[tokio::test] + async fn close_failure_cancellation_and_collision_never_publish_index() { + for failure in ["pack-000000.bin", PACK_INDEX_FILE, "cancel", "collision"] { + let directory = common_test_util::temp_dir::create_temp_dir("packed-failure"); + let store = object_store::secure_fs::SecureFsRoot::open(directory.path()) + .unwrap() + .build_operator(); + if failure == "collision" { + store.write("pack-000000.bin", "winner").await.unwrap(); + } + let factory: MockWriterFactory = Arc::new(move |path, _, writer| { + if path.trim_start_matches('/') == failure { + Box::new(FailClose(writer)) + } else { + writer + } + }); + let store = store.layer( + MockLayerBuilder::default() + .writer_factory(factory) + .build() + .unwrap(), + ); + let shared = PackedWriter::new(store.clone()).unwrap(); + let mut table = PackedTableWriter::new(shared.clone(), "t".into(), 0, false); + table.write(vec![0; 12].into()).await.unwrap(); + table.finish(0, None).await.unwrap(); + let token = CancellationToken::new(); + if failure == "cancel" { + token.cancel(); + } + assert!(shared.lock().await.finish(&token).await.is_err()); + let _ = shared.lock().await.abort().await; + assert!(!store.exists(PACK_INDEX_FILE).await.unwrap()); + if failure == "collision" { + assert_eq!( + store.read("pack-000000.bin").await.unwrap().to_bytes(), + Bytes::from_static(b"winner") + ); + } + } + } +} diff --git a/src/common/datasource/src/parquet_writer.rs b/src/common/datasource/src/parquet_writer.rs index ed35c314339..243a46b4e1a 100644 --- a/src/common/datasource/src/parquet_writer.rs +++ b/src/common/datasource/src/parquet_writer.rs @@ -28,6 +28,7 @@ use tokio_util::sync::CancellationToken; use crate::DEFAULT_WRITE_BUFFER_SIZE; use crate::error::{self, Result}; +use crate::packed_writer::PackedTableWriter; /// Destination creation policy; conditional failures never authorize path deletion. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -52,7 +53,8 @@ pub struct ParquetWriterLimits { /// operation can leave a filesystem worker running after cleanup. pub struct ParquetFileWriter { encoder: Option>>, - sink: Writer, + sink: ParquetSink, + rows: u64, store: ObjectStore, path: String, limits: Option, @@ -60,6 +62,11 @@ pub struct ParquetFileWriter { close_started: bool, } +enum ParquetSink { + Object(Writer), + Packed(Box), +} + impl ParquetFileWriter { /// Open a file using COPY's encoding settings. Destination ownership and /// overwrite policy belong to the caller. None preserves Parquet's defaults. @@ -90,31 +97,7 @@ impl ParquetFileWriter { limits: Option, creation: ParquetCreationPolicy, ) -> Result { - let mut props = WriterProperties::builder() - .set_compression(Compression::ZSTD(ZstdLevel::default())) - .set_statistics_truncate_length(None) - .set_column_index_truncate_length(None); - if let Some(limits) = limits { - ensure!( - limits.row_group_rows > 0 - && limits.flush_threshold_bytes > 0 - && limits.max_row_groups > 0, - error::InvalidParquetWriterLimitsSnafu - ); - props = props - .set_max_row_group_row_count(Some(limits.row_group_rows)) - .set_max_row_group_bytes(None); - } - for field in schema.fields() { - if matches!(field.data_type(), DataType::Timestamp(_, _)) { - let column = ColumnPath::new(vec![field.name().clone()]); - props = props - .set_column_dictionary_enabled(column.clone(), false) - .set_column_encoding(column, Encoding::DELTA_BINARY_PACKED); - } - } - let encoder = ArrowWriter::try_new(Vec::new(), schema, Some(props.build())) - .context(error::WriteParquetSnafu { path })?; + let encoder = build_encoder(schema, limits, path)?; let sink = store .writer_with(path) .concurrent(concurrency) @@ -124,7 +107,8 @@ impl ParquetFileWriter { .context(error::WriteObjectSnafu { path })?; Ok(Self { encoder: Some(encoder), - sink, + sink: ParquetSink::Object(sink), + rows: 0, store, path: path.to_owned(), limits, @@ -133,6 +117,26 @@ impl ParquetFileWriter { }) } + /// Encode into a request-owned pack or standalone destination. + pub fn open_packed( + schema: SchemaRef, + store: ObjectStore, + path: &str, + limits: Option, + packed: PackedTableWriter, + ) -> Result { + Ok(Self { + encoder: Some(build_encoder(schema, limits, path)?), + sink: ParquetSink::Packed(Box::new(packed)), + rows: 0, + store, + path: path.into(), + limits, + creation: ParquetCreationPolicy::IfNotExists, + close_started: false, + }) + } + /// Write a batch, enforcing file limits across batch and row-group boundaries. pub async fn write( &mut self, @@ -182,6 +186,7 @@ impl ParquetFileWriter { check_cancelled(cancellation)?; self.write_bytes(bytes).await?; check_cancelled(cancellation)?; + self.rows += len as u64; offset += len; } Ok(()) @@ -189,11 +194,14 @@ impl ParquetFileWriter { async fn write_bytes(&mut self, bytes: Vec) -> Result<()> { let bytes = Bytes::from(bytes); + let sink = match &mut self.sink { + ParquetSink::Packed(packed) => return packed.write(bytes).await, + ParquetSink::Object(sink) => sink, + }; let chunk = DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize; // Slices retain the complete encoded allocation until its last submission. for offset in (0..bytes.len()).step_by(chunk) { - self.sink - .write(bytes.slice(offset..(offset + chunk).min(bytes.len()))) + sink.write(bytes.slice(offset..(offset + chunk).min(bytes.len()))) .await .context(error::WriteObjectSnafu { path: &self.path })?; } @@ -218,9 +226,12 @@ impl ParquetFileWriter { .context(error::JoinHandleSnafu)??; self.write_bytes(bytes).await?; check_cancelled(cancellation)?; + let sink = match &mut self.sink { + ParquetSink::Packed(packed) => return packed.finish(self.rows, cancellation).await, + ParquetSink::Object(sink) => sink, + }; self.close_started = true; - self.sink - .close() + sink.close() .await .context(error::WriteObjectSnafu { path: &self.path })?; check_cancelled(cancellation)?; @@ -230,7 +241,10 @@ impl ParquetFileWriter { /// Abort after all in-flight operations complete. Preserve ambiguous commits; /// conditional callers delegate cleanup exclusively to the backend. pub async fn abort(mut self) -> Result<()> { - let result = self.sink.abort().await; + let result = match &mut self.sink { + ParquetSink::Packed(packed) => return packed.abort().await, + ParquetSink::Object(sink) => sink.abort().await, + }; if self.creation == ParquetCreationPolicy::Overwrite && result.as_ref().is_err_and(|error| { error.kind() == object_store::ErrorKind::Unsupported @@ -253,6 +267,39 @@ impl ParquetFileWriter { } } +fn build_encoder( + schema: SchemaRef, + limits: Option, + path: &str, +) -> Result>> { + let mut props = WriterProperties::builder() + .set_compression(Compression::ZSTD(ZstdLevel::default())) + .set_statistics_truncate_length(None) + .set_column_index_truncate_length(None); + if let Some(limits) = limits { + ensure!( + limits.row_group_rows > 0 + && limits.flush_threshold_bytes > 0 + && limits.max_row_groups > 0, + error::InvalidParquetWriterLimitsSnafu + ); + props = props + .set_max_row_group_row_count(Some(limits.row_group_rows)) + .set_max_row_group_bytes(None); + } + for field in schema.fields() { + if matches!(field.data_type(), DataType::Timestamp(_, _)) { + let column = ColumnPath::new(vec![field.name().clone()]); + props = props + .set_column_dictionary_enabled(column.clone(), false) + .set_column_encoding(column, Encoding::DELTA_BINARY_PACKED); + } + } + let encoder = ArrowWriter::try_new(Vec::new(), schema, Some(props.build())) + .context(error::WriteParquetSnafu { path })?; + Ok(encoder) +} + fn check_cancelled(cancellation: Option<&CancellationToken>) -> Result<()> { ensure!( cancellation.is_none_or(|token| !token.is_cancelled()), @@ -498,7 +545,7 @@ mod tests { .await .unwrap(); // No storage-layer chunking: observe the application's actual submissions. - writer.sink = store.writer("large.parquet").await.unwrap(); + writer.sink = ParquetSink::Object(store.writer("large.parquet").await.unwrap()); writer.write(batch.clone(), None).await.unwrap(); writer.finish(None).await.unwrap(); let sizes = sizes.lock().unwrap().clone(); @@ -701,7 +748,8 @@ mod tests { .await .unwrap(); // Force footer bytes through the sink before its close operation. - writer.sink = store.writer_with("cancel.parquet").chunk(1).await.unwrap(); + writer.sink = + ParquetSink::Object(store.writer_with("cancel.parquet").chunk(1).await.unwrap()); let cancellation = CancellationToken::new(); let result = { let finish = writer.finish(Some(&cancellation)); diff --git a/src/operator/src/statement/copy_table_to.rs b/src/operator/src/statement/copy_table_to.rs index d3d44f7655f..f9727be6c12 100644 --- a/src/operator/src/statement/copy_table_to.rs +++ b/src/operator/src/statement/copy_table_to.rs @@ -24,6 +24,7 @@ use common_datasource::file_format::csv::stream_to_csv; use common_datasource::file_format::json::stream_to_json; use common_datasource::file_format::parquet::stream_to_parquet; use common_datasource::object_store::build_backend_for_write_with_path; +use common_datasource::packed_writer::PackedTableWriter; use common_datasource::parquet_writer::ParquetFileWriter; use common_query::Output; use common_recordbatch::adapter::DfRecordBatchStreamAdapter; @@ -127,7 +128,7 @@ impl StatementExecutor { req: CopyTableRequest, query_ctx: QueryContextRef, ) -> Result { - self.copy_captured_table_to_managed(table, req, query_ctx, None) + self.copy_captured_table_to_managed(table, req, query_ctx, None, None) .await } @@ -137,6 +138,7 @@ impl StatementExecutor { req: CopyTableRequest, query_ctx: QueryContextRef, managed: Option<(&ExportWriteBudget, &CancellationToken)>, + packed: Option, ) -> Result { let info = table.table_info(); let table_ref = TableReference::full(&info.catalog_name, &info.schema_name, &info.name); @@ -185,7 +187,7 @@ impl StatementExecutor { } = &req; debug!("Copy table: {table_id} to location: {location}"); - self.copy_to_file_managed(&format, output, location, connection, managed) + self.copy_to_file_managed(&format, output, location, connection, managed, packed) .await } @@ -196,7 +198,7 @@ impl StatementExecutor { location: &str, connection: &HashMap, ) -> Result { - self.copy_to_file_managed(format, output, location, connection, None) + self.copy_to_file_managed(format, output, location, connection, None, None) .await } @@ -207,6 +209,7 @@ impl StatementExecutor { location: &str, connection: &HashMap, managed: Option<(&ExportWriteBudget, &CancellationToken)>, + packed: Option, ) -> Result { let output = if managed.is_none() { output @@ -229,7 +232,15 @@ impl StatementExecutor { violated: format!("Expected filename, path: {location}"), })?; if let Some((budget, token)) = managed { - stream_to_managed_parquet(stream, backend.object_store, &filename, budget, token).await + stream_to_managed_parquet_with_packed( + stream, + backend.object_store, + &filename, + budget, + token, + packed, + ) + .await } else { self.stream_to_file(stream, format, backend.object_store, &filename) .await @@ -237,12 +248,24 @@ impl StatementExecutor { } } +#[cfg(test)] pub(crate) async fn stream_to_managed_parquet( + stream: SendableRecordBatchStream, + store: ObjectStore, + path: &str, + budget: &ExportWriteBudget, + token: &CancellationToken, +) -> Result { + stream_to_managed_parquet_with_packed(stream, store, path, budget, token, None).await +} + +async fn stream_to_managed_parquet_with_packed( mut stream: SendableRecordBatchStream, store: ObjectStore, path: &str, budget: &ExportWriteBudget, token: &CancellationToken, + packed: Option, ) -> Result { use common_recordbatch::{RecordBatch, map_dictionary_to_values_schema}; let original = stream.schema(); @@ -261,10 +284,21 @@ pub(crate) async fn stream_to_managed_parquet( } else { expanded_schema.clone() }; - let mut writer = - ParquetFileWriter::open(output_schema.arrow_schema().clone(), store, path, 1, None) - .await - .context(error::WriteStreamToFileSnafu { path })?; + let packed_destination = packed.is_some(); + let mut writer = if let Some(packed) = packed { + ParquetFileWriter::open_packed( + output_schema.arrow_schema().clone(), + store, + path, + Some( + crate::statement::export_logical_tables::LogicalTableExportLimits::default().writer, + ), + packed, + ) + } else { + ParquetFileWriter::open(output_schema.arrow_schema().clone(), store, path, 1, None).await + } + .context(error::WriteStreamToFileSnafu { path })?; let mut started = false; let result = async { let mut rows = 0; @@ -358,7 +392,9 @@ pub(crate) async fn stream_to_managed_parquet( .await; if result.is_err() { token.cancel(); - if started && let Err(error) = writer.abort().await { + if (started || packed_destination) + && let Err(error) = writer.abort().await + { common_telemetry::warn!(error; "Failed to abort ordinary export file"); } } diff --git a/src/operator/src/statement/database_copy.rs b/src/operator/src/statement/database_copy.rs index 9878560afaf..08162c5d5e6 100644 --- a/src/operator/src/statement/database_copy.rs +++ b/src/operator/src/statement/database_copy.rs @@ -60,7 +60,7 @@ pub(crate) fn parse_parallelism_from_option_map(options: &HashMap) -> Result<()> { if let Some(layout) = options.get("metric_data_layout") { return error::InvalidCopyParameterSnafu { diff --git a/src/operator/src/statement/export_database.rs b/src/operator/src/statement/export_database.rs index dbe18ba763a..f363c73d7a2 100644 --- a/src/operator/src/statement/export_database.rs +++ b/src/operator/src/statement/export_database.rs @@ -20,6 +20,7 @@ use std::future::Future; use common_datasource::file_format::Format; use common_datasource::object_store::{FILE_SCHEMA, FS_SCHEMA, build_backend_for_write, parse_url}; +use common_datasource::packed_writer::{PackedTableWriter, PackedWriter}; use common_meta::key::table_route::TableRouteValue; use futures::StreamExt; use futures::stream::FuturesUnordered; @@ -70,7 +71,13 @@ impl StatementExecutor { req: CopyDatabaseRequest, tables: Vec, ) -> Result { - validate_database_export_layout(&req.with)?; + if req + .with + .get("metric_data_layout") + .is_none_or(|v| v != "packed") + { + validate_database_export_layout(&req.with)?; + } validate_database_directory(&req.location)?; let format = Format::try_from(&req.with).context(error::ParseFileFormatSnafu)?; ensure!( @@ -205,9 +212,27 @@ impl StatementExecutor { let req = &plan.request; let parallelism = parse_parallelism_from_option_map(&req.with); let budget = ExportWriteBudget::new(parallelism); + let packed = if req + .with + .get("metric_data_layout") + .is_some_and(|v| v == "packed") + { + let store = + build_backend_for_write(&req.location, &req.connection, &self.local_file_access) + .await + .context(error::BuildBackendSnafu)?; + Some( + PackedWriter::new(store).context(error::WriteStreamToFileSnafu { + path: &req.location, + })?, + ) + } else { + None + }; let rows = run_database_export_jobs(plan.jobs, parallelism, cancellation, |job, token| { let ctx = ctx.clone(); let budget = budget.clone(); + let packed = packed.clone(); async move { match job { DatabaseExportJob::Metric(unit) => self @@ -220,6 +245,7 @@ impl StatementExecutor { &token, ctx, budget, + packed, ) .await .map(|summary| summary.rows), @@ -238,19 +264,57 @@ impl StatementExecutor { limit: None, }; let _permit = budget.writer(&token).await?; + let destination = packed.map(|shared| { + PackedTableWriter::new(shared, info.name.clone(), info.table_id(), true) + }); self.copy_captured_table_to_managed( table, copy, ctx, Some((&budget, &token)), + destination, ) .await } } } }) - .await?; - Ok(DatabaseExportSummary { rows, output_files }) + .await; + if let Some(packed) = packed { + let mut packed = packed.lock().await; + let result = match rows { + Ok(rows) => packed + .finish(cancellation) + .await + .context(error::WriteStreamToFileSnafu { + path: &req.location, + }) + .and_then(|files| { + Ok(DatabaseExportSummary { + rows, + output_files: files + .into_iter() + .map(|file| { + DatabaseExportFile::new(&req.location, &file, "") + .map(|f| f.location) + }) + .collect::>()?, + }) + }), + Err(error) => Err(error), + }; + if result.is_err() + && let Err(error) = packed.abort().await + { + common_telemetry::warn!(error; "Failed to abort Metric pack"); + } + result + } else { + Ok(DatabaseExportSummary { + rows: rows?, + output_files, + }) + } } } diff --git a/src/operator/src/statement/export_logical_tables.rs b/src/operator/src/statement/export_logical_tables.rs index 4a2e186d477..368335f8450 100644 --- a/src/operator/src/statement/export_logical_tables.rs +++ b/src/operator/src/statement/export_logical_tables.rs @@ -31,6 +31,7 @@ use arrow::datatypes::{DataType, SchemaRef}; use arrow::downcast_dictionary_array; use arrow::record_batch::RecordBatch; use common_datasource::object_store::build_backend_for_write; +use common_datasource::packed_writer::{PackedTableWriter, PackedWriterRef}; use common_datasource::parquet_writer::{ ParquetCreationPolicy, ParquetFileWriter, ParquetWriterLimits, }; @@ -108,6 +109,7 @@ pub struct LogicalTableExport { } pub(crate) struct LogicalTableProjection { + name: String, output: DatabaseExportFile, schema: SchemaRef, projection: Vec, @@ -196,6 +198,7 @@ impl LogicalTableExport { .insert( info.table_id(), LogicalTableProjection { + name: name.clone(), output: DatabaseExportFile::new(directory, name, ".parquet")?, schema, projection: indices, @@ -322,6 +325,7 @@ impl StatementExecutor { &cancellation.child_token(), query_ctx, ExportWriteBudget::new(1), + None, ) .await } @@ -337,6 +341,7 @@ impl StatementExecutor { cancellation: &CancellationToken, query_ctx: QueryContextRef, budget: Arc, + packed: Option, ) -> Result { limits.validate()?; let (store, stream) = tokio::select! { @@ -356,7 +361,7 @@ impl StatementExecutor { Ok((store, stream)) } => result?, }; - export_stream_managed(unit, stream, &store, limits, cancellation, budget).await + export_stream_managed(unit, stream, &store, limits, cancellation, budget, packed).await } } @@ -375,6 +380,7 @@ async fn export_stream( limits, &cancellation.child_token(), ExportWriteBudget::new(1), + None, ) .await } @@ -386,8 +392,10 @@ async fn export_stream_managed( limits: LogicalTableExportLimits, cancellation: &CancellationToken, budget: Arc, + packed: Option, ) -> Result { let mut writers = TableWriters::new(budget.clone()); + writers.packed = packed; let result = write_tables( unit, stream, @@ -535,6 +543,26 @@ struct ActiveWriter { } impl ActiveWriter { + fn open_packed( + table: &LogicalTableProjection, + id: u32, + store: &ObjectStore, + limits: LogicalTableExportLimits, + shared: PackedWriterRef, + ) -> Result { + let path = format!("table-{id}.parquet"); + let packed = PackedTableWriter::new(shared, table.name.clone(), id, false); + let writer = ParquetFileWriter::open_packed( + table.schema.clone(), + store.clone(), + &path, + Some(limits.writer), + packed, + ) + .map_err(|e| map_writer_error(e, &path))?; + Ok(Self { path, writer }) + } + async fn open( table: &LogicalTableProjection, store: &ObjectStore, diff --git a/src/operator/src/statement/export_logical_tables/tests.rs b/src/operator/src/statement/export_logical_tables/tests.rs index a03fdf1a80a..7aa6daab127 100644 --- a/src/operator/src/statement/export_logical_tables/tests.rs +++ b/src/operator/src/statement/export_logical_tables/tests.rs @@ -146,6 +146,7 @@ async fn routes_across_batches_and_writes_empty_files() { export_limits(), &CancellationToken::new(), budget.clone(), + None, ) .await .unwrap(); @@ -400,12 +401,17 @@ async fn native_histogram_parquet_roundtrip() { struct PausedFileWriter { inner: Option, + paused: bool, started: Arc, release: Arc, } impl object_store::layers::mock::oio::Write for PausedFileWriter { async fn write(&mut self, bytes: object_store::Buffer) -> object_store::Result<()> { + if self.paused { + return self.inner.as_mut().unwrap().write(bytes).await; + } + self.paused = true; let mut inner = self.inner.take().unwrap(); let started = self.started.clone(); let release = self.release.clone(); @@ -439,7 +445,7 @@ impl object_store::layers::mock::oio::Write for PausedFileWriter { #[tokio::test] async fn cancellation_drains_storage_and_preserves_committed_files() { - for large in [false, true] { + for (large, packed) in [(false, false), (true, false), (false, true)] { use object_store::layers::mock::{MockLayerBuilder, MockWriterFactory}; let directory = common_test_util::temp_dir::create_temp_dir("metric_export_pending_open"); let access = @@ -459,6 +465,7 @@ async fn cancellation_drains_storage_and_preserves_committed_files() { move |_, _, inner| { Box::new(PausedFileWriter { inner: Some(inner), + paused: false, started: started.clone(), release: release.clone(), }) @@ -470,6 +477,16 @@ async fn cancellation_drains_storage_and_preserves_committed_files() { .build() .unwrap(), ); + let store = store.layer(object_store::layers::CapabilityOverrideLayer::new( + move |mut capability| { + if packed { + capability.write_multi_max_size = Some(256); + } + capability + }, + )); + let packed = packed + .then(|| common_datasource::packed_writer::PackedWriter::new(store.clone()).unwrap()); let cancellation = CancellationToken::new(); let (unit, input) = if large { let field = Field::new("value", DataType::Utf8, true); @@ -518,6 +535,7 @@ async fn cancellation_drains_storage_and_preserves_committed_files() { limits, &cancellation, budget.clone(), + packed.clone(), ); tokio::pin!(export); tokio::select! { @@ -528,6 +546,7 @@ async fn cancellation_drains_storage_and_preserves_committed_files() { assert!(budget.available().1 < 64 * 1024 * 1024); } let held = budget.available(); + assert_eq!(held.0, 0); cancellation.cancel(); assert!(futures::poll!(&mut export).is_pending()); assert_eq!(budget.available().0, held.0); @@ -542,8 +561,16 @@ async fn cancellation_drains_storage_and_preserves_committed_files() { result, Err(error::Error::LogicalTableExportCancelled { .. }) )); - assert_eq!(store.exists("cpu.v1.parquet").await.unwrap(), !large); - if !large { + if let Some(packed) = &packed { + packed.lock().await.abort().await.unwrap(); + assert!(!store.exists("pack-000000.bin").await.unwrap()); + assert!(!store.exists("pack-index.json").await.unwrap()); + } + assert_eq!( + store.exists("cpu.v1.parquet").await.unwrap(), + !large && packed.is_none() + ); + if !large && packed.is_none() { let (_, batches) = read(&store, "cpu.v1.parquet").await; assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 1); } @@ -829,7 +856,8 @@ async fn groups_and_ordinary_files_share_writer_admission_and_drain() { &store, export_limits(), &token, - budget.clone() + budget.clone(), + None, ), export_stream_managed( &b, @@ -837,7 +865,8 @@ async fn groups_and_ordinary_files_share_writer_admission_and_drain() { &store, export_limits(), &token, - budget.clone() + budget.clone(), + None, ), ordinary, ) diff --git a/src/operator/src/statement/export_logical_tables/writers.rs b/src/operator/src/statement/export_logical_tables/writers.rs index 38a63936a5b..70b3bcd2e32 100644 --- a/src/operator/src/statement/export_logical_tables/writers.rs +++ b/src/operator/src/statement/export_logical_tables/writers.rs @@ -15,6 +15,7 @@ use std::sync::Arc; use arrow::record_batch::RecordBatch; +use common_datasource::packed_writer::PackedWriterRef; use futures::stream::FuturesUnordered; use futures::{FutureExt, StreamExt}; use object_store::ObjectStore; @@ -111,6 +112,7 @@ pub(crate) struct Payload { } pub(crate) struct TableWriters { + pub(crate) packed: Option, current: Option<(u32, mpsc::Sender)>, tasks: FuturesUnordered>>, budget: Arc, @@ -120,6 +122,7 @@ impl TableWriters { pub(crate) fn new(budget: Arc) -> Self { Self { current: None, + packed: None, tasks: FuturesUnordered::new(), budget, } @@ -150,7 +153,11 @@ impl TableWriters { self.close_input(); let permit = self.budget.writer(token).await?; self.reap_for_admission(token).await?; - let writer = ActiveWriter::open(table, store, limits).await?; + let writer = if let Some(packed) = &self.packed { + ActiveWriter::open_packed(table, id, store, limits, packed.clone())? + } else { + ActiveWriter::open(table, store, limits).await? + }; let (sender, receiver) = mpsc::channel(2); self.current = Some((id, sender)); let token = token.clone(); diff --git a/src/servers/src/http.rs b/src/servers/src/http.rs index a5266576cbe..b9fcd0a9420 100644 --- a/src/servers/src/http.rs +++ b/src/servers/src/http.rs @@ -1381,7 +1381,9 @@ impl HttpServer { .route( "/capabilities", routing::get(|| async { - axum::Json(serde_json::json!({"metric_packed_import": 1})) + axum::Json( + serde_json::json!({"metric_packed_import": 1, "metric_packed_export": 1}), + ) }), ) .route( @@ -1897,7 +1899,7 @@ mod test { assert_eq!(response.status(), StatusCode::OK); assert_eq!( response.json::().await, - serde_json::json!({"metric_packed_import": 1}) + serde_json::json!({"metric_packed_import": 1, "metric_packed_export": 1}) ); } } diff --git a/tests-integration/tests/export_logical_tables.rs b/tests-integration/tests/export_logical_tables.rs index 1f9161d970f..80169a391cd 100644 --- a/tests-integration/tests/export_logical_tables.rs +++ b/tests-integration/tests/export_logical_tables.rs @@ -1451,7 +1451,7 @@ impl servers::interceptor::SqlQueryInterceptor for FailSecondChunk { #[tokio::test] async fn metric_export_v2_cli_resume_roundtrip() { - metric_export_v2_cli_roundtrip(false).await; + metric_export_v2_cli_roundtrip(false, false).await; } #[tokio::test] @@ -1461,11 +1461,54 @@ async fn metric_export_v2_cli_s3_resume_roundtrip() { if std::env::var("GT_S3_ENDPOINT_URL").is_ok_and(|e| !e.is_empty()) && std::env::var("GT_S3_BUCKET").is_ok_and(|b| !b.is_empty()) { - metric_export_v2_cli_roundtrip(true).await; + metric_export_v2_cli_roundtrip(true, false).await; } } -async fn metric_export_v2_cli_roundtrip(s3: bool) { +#[tokio::test(flavor = "multi_thread")] +async fn packed_export_v2_cli_local_roundtrip() { + metric_export_v2_cli_roundtrip(false, true).await; +} + +#[tokio::test(flavor = "multi_thread")] +async fn packed_export_v2_cli_s3_roundtrip() { + if std::env::var("GT_S3_ENDPOINT_URL").is_ok_and(|e| !e.is_empty()) + && std::env::var("GT_S3_BUCKET").is_ok_and(|b| !b.is_empty()) + { + metric_export_v2_cli_roundtrip(true, true).await; + } +} + +async fn exported_table_path( + store: &object_store::ObjectStore, + packed: bool, + chunk: u32, + name: &str, +) -> String { + let prefix = format!("data/public/{chunk}/"); + if !packed { + return format!("{prefix}{name}.parquet"); + } + let index: common_datasource::packed_snapshot::PackIndex = serde_json::from_slice( + &store + .read(&format!("{prefix}pack-index.json")) + .await + .unwrap() + .to_bytes(), + ) + .unwrap(); + format!( + "{prefix}{}", + index + .tables + .iter() + .find(|t| t.table_name == name) + .unwrap() + .object + ) +} + +async fn metric_export_v2_cli_roundtrip(s3: bool, packed: bool) { use servers::interceptor::SqlQueryInterceptorRef; let plugins = common_base::Plugins::new(); let faults = Arc::new(FailSecondChunk { @@ -1499,6 +1542,9 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { sql(instance, "INSERT INTO audit VALUES ('a',1,1),('z',NULL,3)").await; sql(instance, "CREATE VIEW dashboard AS SELECT * FROM audit").await; sql(instance, "CREATE DATABASE z_later").await; + if packed { + sql(instance, "CREATE DATABASE empty_schema").await; + } sql( instance, "CREATE TABLE z_later.events (ts TIMESTAMP TIME INDEX, val BIGINT)", @@ -1530,6 +1576,9 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap(); for experimental in [false, true] { for layout in ["packed", "invalid"] { + if experimental && layout == "packed" { + continue; + } let path = destination .path() .join(format!("unsupported-{experimental}-{layout}")); @@ -1544,7 +1593,7 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { ) .await .remove(0); - let error = result.expect_err("export must reject import-only layouts"); + let error = result.expect_err("export must reject unsupported layouts"); assert!( format!("{error:?}").contains("metric_data_layout"), "{error:?}" @@ -1552,6 +1601,41 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { assert!(!path.exists()); } } + if packed && !s3 { + for window in [None, Some("3ms")] { + let path = destination.path().join("empty-range"); + let uri = url::Url::from_directory_path(&path).unwrap(); + let mut args = vec![ + "export-v2", + "create", + "--addr", + &addr, + "--to", + uri.as_str(), + "--schemas", + "public", + "--experimental-metric-export", + "--metric-data-layout", + "packed", + "--no-proxy", + "--start-time", + "1970-01-01T00:00:00Z", + "--end-time", + "1970-01-01T00:00:00Z", + ]; + if let Some(window) = window { + args.extend(["--chunk-time-window", window]); + } + let error = run_data_cli(&args).await.unwrap_err(); + assert!( + error + .to_string() + .contains("Packed export requires --start-time to be earlier than --end-time"), + "{error}" + ); + assert!(!path.exists()); + } + } let (uri, store, storage_args) = if s3 { let endpoint = std::env::var("GT_S3_ENDPOINT_URL").unwrap(); let bucket = std::env::var("GT_S3_BUCKET").unwrap(); @@ -1604,7 +1688,11 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { "--to", &uri, "--schemas", - "public,z_later", + if packed { + "public,z_later,empty_schema" + } else { + "public,z_later" + }, "--experimental-metric-export", "--no-proxy", "--start-time", @@ -1617,6 +1705,9 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { "never", ]; args.extend(storage_args.iter().map(String::as_str)); + if packed { + args.extend(["--metric-data-layout", "packed"]); + } assert!(run_data_cli(&args).await.is_err()); let before: cli::export_v2::manifest::Manifest = serde_json::from_slice(&store.read("manifest.json").await.unwrap().to_vec()).unwrap(); @@ -1631,7 +1722,7 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { if s3 { assert!( store - .stat("data/public/1/bulk.parquet") + .stat(&exported_table_path(&store, packed, 1, "bulk").await) .await .unwrap() .content_length() @@ -1642,7 +1733,12 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { for path in &before.chunks[0].files { preserved.push((path.clone(), store.read(path).await.unwrap().to_vec())); } - assert!(store.exists("data/public/2/audit.parquet").await.unwrap()); + assert!( + store + .exists(&exported_table_path(&store, packed, 2, "audit").await) + .await + .unwrap() + ); store .write("data/public/2/unknown.txt", "keep") .await @@ -1666,6 +1762,59 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { let after: cli::export_v2::manifest::Manifest = serde_json::from_slice(&store.read("manifest.json").await.unwrap().to_vec()).unwrap(); assert!(after.is_complete()); + assert_eq!(after.version, if packed { 2 } else { 1 }); + assert_eq!(before.data_layout, after.data_layout); + if packed { + use common_datasource::packed_snapshot::{ObjectKind, PackIndex}; + for chunk in &after.chunks { + let empty: PackIndex = serde_json::from_slice( + &store + .read(&format!("data/empty_schema/{}/pack-index.json", chunk.id)) + .await + .unwrap() + .to_bytes(), + ) + .unwrap(); + empty.validate_membership([]).unwrap(); + let prefix = format!("data/public/{}/", chunk.id); + let index: PackIndex = serde_json::from_slice( + &store + .read(&format!("{prefix}pack-index.json")) + .await + .unwrap() + .to_bytes(), + ) + .unwrap(); + index + .validate_membership(names.iter().map(String::as_str)) + .unwrap(); + assert_eq!( + index + .objects + .iter() + .filter(|o| o.kind == ObjectKind::Pack) + .count(), + 1 + ); + assert!(index.tables.iter().filter(|t| t.row_count == 0).count() >= 2); + let mut expected = index + .objects + .iter() + .map(|o| format!("{prefix}{}", o.path)) + .collect::>(); + expected.push(format!("{prefix}pack-index.json")); + expected.sort(); + assert_eq!( + chunk + .files + .iter() + .filter(|f| f.starts_with(&prefix)) + .cloned() + .collect::>(), + expected + ); + } + } assert_eq!( serde_json::to_value(&before.chunks[0]).unwrap(), serde_json::to_value(&after.chunks[0]).unwrap() @@ -1727,6 +1876,39 @@ async fn metric_export_v2_cli_roundtrip(s3: bool) { ) .await ); + if packed { + let schema_uri = format!("{uri}/schema-only"); + let mut schema_args = vec![ + "export-v2", + "create", + "--addr", + &addr, + "--to", + &schema_uri, + "--schemas", + "public", + "--schema-only", + "--experimental-metric-export", + "--metric-data-layout", + "packed", + "--no-proxy", + "--progress", + "never", + ]; + schema_args.extend(storage_args.iter().map(String::as_str)); + run_data_cli(&schema_args).await.unwrap(); + let schema_manifest: cli::export_v2::manifest::Manifest = serde_json::from_slice( + &store + .read("schema-only/manifest.json") + .await + .unwrap() + .to_bytes(), + ) + .unwrap(); + assert!(schema_manifest.schema_only && schema_manifest.chunks.is_empty()); + assert_eq!(schema_manifest.version, 1); + assert!(schema_manifest.data_layout.is_none()); + } server.abort(); target_server.abort(); if s3 { diff --git a/tests/compatibility/cases/snapshot_parquet_copy/setup.sql b/tests/compatibility/cases/snapshot_parquet_copy/setup.sql index 8342e1847b5..38bc7c74aec 100644 --- a/tests/compatibility/cases/snapshot_parquet_copy/setup.sql +++ b/tests/compatibility/cases/snapshot_parquet_copy/setup.sql @@ -1,4 +1,6 @@ -CREATE TABLE snapshot_values (ts TIMESTAMP TIME INDEX, val BIGINT); -INSERT INTO snapshot_values VALUES (1, 42), (2, 7); +CREATE TABLE snapshot_values (ts TIMESTAMP TIME INDEX, val BIGINT, host STRING); +INSERT INTO snapshot_values VALUES (1, 42, ''), (2, 7, NULL); COPY snapshot_values TO '${SQLNESS_HOME}/snapshot_parquet_copy/values.parquet' WITH (FORMAT='parquet'); TRUNCATE TABLE snapshot_values; +CREATE TABLE snapshot_empty (ts TIMESTAMP TIME INDEX, val DOUBLE); +COPY snapshot_empty TO '${SQLNESS_HOME}/snapshot_parquet_copy/empty.parquet' WITH (FORMAT='parquet'); diff --git a/tests/compatibility/cases/snapshot_parquet_copy/verify.result b/tests/compatibility/cases/snapshot_parquet_copy/verify.result index d3e4d2483ab..6c590e1a1ad 100644 --- a/tests/compatibility/cases/snapshot_parquet_copy/verify.result +++ b/tests/compatibility/cases/snapshot_parquet_copy/verify.result @@ -4,9 +4,30 @@ Affected Rows: 2 SELECT * FROM snapshot_values ORDER BY ts; -+-------------------------+-----+ -| ts | val | -+-------------------------+-----+ -| 1970-01-01T00:00:00.001 | 42 | -| 1970-01-01T00:00:00.002 | 7 | -+-------------------------+-----+ ++-------------------------+-----+------+ +| ts | val | host | ++-------------------------+-----+------+ +| 1970-01-01T00:00:00.001 | 42 | | +| 1970-01-01T00:00:00.002 | 7 | | ++-------------------------+-----+------+ + +SELECT val, host IS NULL AS host_is_null, host = '' AS host_is_empty FROM snapshot_values ORDER BY ts; + ++-----+--------------+---------------+ +| val | host_is_null | host_is_empty | ++-----+--------------+---------------+ +| 42 | false | true | +| 7 | true | | ++-----+--------------+---------------+ + +COPY snapshot_empty FROM '${SQLNESS_HOME}/snapshot_parquet_copy/empty.parquet' WITH (FORMAT='parquet'); + +Affected Rows: 0 + +SELECT COUNT(*) FROM snapshot_empty; + ++----------+ +| count(*) | ++----------+ +| 0 | ++----------+ diff --git a/tests/compatibility/cases/snapshot_parquet_copy/verify.sql b/tests/compatibility/cases/snapshot_parquet_copy/verify.sql index c782dcacb1b..0e35e3c127f 100644 --- a/tests/compatibility/cases/snapshot_parquet_copy/verify.sql +++ b/tests/compatibility/cases/snapshot_parquet_copy/verify.sql @@ -1,2 +1,5 @@ COPY snapshot_values FROM '${SQLNESS_HOME}/snapshot_parquet_copy/values.parquet' WITH (FORMAT='parquet'); SELECT * FROM snapshot_values ORDER BY ts; +SELECT val, host IS NULL AS host_is_null, host = '' AS host_is_empty FROM snapshot_values ORDER BY ts; +COPY snapshot_empty FROM '${SQLNESS_HOME}/snapshot_parquet_copy/empty.parquet' WITH (FORMAT='parquet'); +SELECT COUNT(*) FROM snapshot_empty;