feat: export Metric snapshots with packed Parquet objects (#9382)

* feat: export Metric snapshots with packed Parquet objects

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: reject empty packed export time ranges

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test: cover cancellation during packed export I/O

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
This commit is contained in:
jeremyhi
2026-09-28 15:22:02 +00:00
committed by GitHub
parent e2b7e5e0a8
commit 3c0e2a8d55
24 changed files with 1584 additions and 127 deletions
Generated
+2
View File
@@ -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",
+1
View File
@@ -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
+3 -1
View File
@@ -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;
+65 -4
View File
@@ -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<String>,
/// 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<Duration>,
@@ -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,
+189 -1
View File
@@ -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<object_store::layers::mock::oio::Writer>,
started: std::sync::Arc<tokio::sync::Notify>,
release: std::sync::Arc<tokio::sync::Notify>,
}
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<object_store::layers::mock::Metadata> {
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);
+4
View File
@@ -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='{}'",
+1 -1
View File
@@ -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.
+93 -9
View File
@@ -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());
+62 -40
View File
@@ -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();
}
}
}
+1
View File
@@ -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
+1
View File
@@ -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)]
+669
View File
@@ -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<Mutex<PackedWriter>>;
#[derive(Clone, Copy)]
struct UploadLimits {
part: usize,
object: u64,
}
impl UploadLimits {
fn for_store(store: &ObjectStore) -> Result<Self> {
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<Self> {
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<T: serde::Serialize>(
&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<Upload>,
index: PackIndex,
next_pack: usize,
}
impl PackedWriter {
pub fn new(store: ObjectStore) -> Result<PackedWriterRef> {
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<Vec<String>> {
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<u8>,
ordinary: bool,
standalone: Option<Upload>,
}
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::<Vec<_>>();
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::<Vec<_>>()
.await
.unwrap()
} else {
ParquetRecordBatchReaderBuilder::try_new(
store.read(&object.path).await.unwrap().to_bytes(),
)
.unwrap()
.build()
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.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<std::sync::Mutex<Vec<usize>>>);
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<Metadata> {
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<Metadata> {
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")
);
}
}
}
}
+82 -34
View File
@@ -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<ArrowWriter<Vec<u8>>>,
sink: Writer,
sink: ParquetSink,
rows: u64,
store: ObjectStore,
path: String,
limits: Option<ParquetWriterLimits>,
@@ -60,6 +62,11 @@ pub struct ParquetFileWriter {
close_started: bool,
}
enum ParquetSink {
Object(Writer),
Packed(Box<PackedTableWriter>),
}
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<ParquetWriterLimits>,
creation: ParquetCreationPolicy,
) -> Result<Self> {
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<ParquetWriterLimits>,
packed: PackedTableWriter,
) -> Result<Self> {
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<u8>) -> 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<ParquetWriterLimits>,
path: &str,
) -> Result<ArrowWriter<Vec<u8>>> {
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));
+45 -9
View File
@@ -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<usize> {
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<PackedTableWriter>,
) -> Result<usize> {
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<String, String>,
) -> Result<usize> {
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<String, String>,
managed: Option<(&ExportWriteBudget, &CancellationToken)>,
packed: Option<PackedTableWriter>,
) -> Result<usize> {
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<usize> {
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<PackedTableWriter>,
) -> Result<usize> {
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");
}
}
+1 -1
View File
@@ -60,7 +60,7 @@ pub(crate) fn parse_parallelism_from_option_map(options: &HashMap<String, String
.clamp(1, Semaphore::MAX_PERMITS)
}
/// Rejects import-only layouts before either database export path creates output.
/// Rejects layouts unsupported by the per-table export path.
pub(crate) fn validate_database_export_layout(options: &HashMap<String, String>) -> Result<()> {
if let Some(layout) = options.get("metric_data_layout") {
return error::InvalidCopyParameterSnafu {
+67 -3
View File
@@ -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<TableRef>,
) -> Result<PreparedDatabaseExport> {
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::<Result<_>>()?,
})
}),
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,
})
}
}
}
@@ -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<usize>,
@@ -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<ExportWriteBudget>,
packed: Option<PackedWriterRef>,
) -> Result<LogicalTableExportSummary> {
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<ExportWriteBudget>,
packed: Option<PackedWriterRef>,
) -> Result<LogicalTableExportSummary> {
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<Self> {
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,
@@ -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<object_store::layers::mock::oio::Writer>,
paused: bool,
started: Arc<tokio::sync::Notify>,
release: Arc<tokio::sync::Notify>,
}
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::<usize>(), 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,
)
@@ -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<PackedWriterRef>,
current: Option<(u32, mpsc::Sender<Payload>)>,
tasks: FuturesUnordered<JoinHandle<Result<()>>>,
budget: Arc<ExportWriteBudget>,
@@ -120,6 +122,7 @@ impl TableWriters {
pub(crate) fn new(budget: Arc<ExportWriteBudget>) -> 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();
+4 -2
View File
@@ -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::<serde_json::Value>().await,
serde_json::json!({"metric_packed_import": 1})
serde_json::json!({"metric_packed_import": 1, "metric_packed_export": 1})
);
}
}
@@ -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::<Vec<_>>();
expected.push(format!("{prefix}pack-index.json"));
expected.sort();
assert_eq!(
chunk
.files
.iter()
.filter(|f| f.starts_with(&prefix))
.cloned()
.collect::<Vec<_>>(),
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 {
@@ -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');
@@ -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 |
+----------+
@@ -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;