From 2066d67aa69a0e62e73024c056eb5ddf05305c5e Mon Sep 17 00:00:00 2001 From: Sonic Shih Date: Mon, 3 Aug 2026 18:24:58 +0800 Subject: [PATCH] fix(collector): speed up Polymarket tape upload with zstd threads and ossutil multipart tuning Tick-level tapes (20-25 GB) were compressed with 'zstd -T1' and copied without ossutil multipart flags, capping uploads at ~50 MiB/s. - zstd thread count is configurable via ZSTD_THREADS (default 0 = auto, all cores); compression level and timeout semantics unchanged. - ossutil cp gains --parallel (OSS_PARALLEL, default 8) and --part-size (OSS_PART_SIZE, default 32Mi) on every cp invocation in the upload pipeline: upload, pre-upload existence-check download, and readback verify download. - Env values parse fail-closed through the existing env_u64/env_or helpers, and UploadConfig::validate rejects zero parallelism and an empty part size. Refs #655 --- .../collector/src/bin/polymarket-raw-ops.rs | 6 + .../tools/collector/src/polymarket_upload.rs | 119 +++++++++++++++++- 2 files changed, 120 insertions(+), 5 deletions(-) diff --git a/rust_hft/tools/collector/src/bin/polymarket-raw-ops.rs b/rust_hft/tools/collector/src/bin/polymarket-raw-ops.rs index 67f42844d..106470493 100644 --- a/rust_hft/tools/collector/src/bin/polymarket-raw-ops.rs +++ b/rust_hft/tools/collector/src/bin/polymarket-raw-ops.rs @@ -275,6 +275,9 @@ async fn run(cli: Cli) -> Result<()> { } => { let zstd_timeout = env_u64(zstd_timeout, "ZSTD_TIMEOUT_SECONDS", 300)?; let oss_timeout = env_u64(oss_timeout, "OSS_COPY_TIMEOUT_SECONDS", 300)?; + let zstd_threads = env_u64(None, "ZSTD_THREADS", 0)?; + let oss_parallel = env_u64(None, "OSS_PARALLEL", 8)?; + let oss_part_size = env_or(None, "OSS_PART_SIZE", "32Mi"); let max_concurrent_uploads = upload_concurrency.unwrap_or(DEFAULT_MAX_CONCURRENT_UPLOADS); let config = UploadConfig { @@ -293,6 +296,9 @@ async fn run(cli: Cli) -> Result<()> { zstd_timeout: Duration::from_secs(zstd_timeout), oss_timeout: Duration::from_secs(oss_timeout), max_concurrent_uploads, + zstd_threads, + oss_parallel, + oss_part_size, }; println!( "{}", diff --git a/rust_hft/tools/collector/src/polymarket_upload.rs b/rust_hft/tools/collector/src/polymarket_upload.rs index d256e1bc0..fdb98d1d8 100644 --- a/rust_hft/tools/collector/src/polymarket_upload.rs +++ b/rust_hft/tools/collector/src/polymarket_upload.rs @@ -94,6 +94,9 @@ pub struct UploadConfig { pub zstd_timeout: Duration, pub oss_timeout: Duration, pub max_concurrent_uploads: usize, + pub zstd_threads: u64, + pub oss_parallel: u64, + pub oss_part_size: String, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] @@ -123,6 +126,12 @@ impl UploadConfig { if !(1..=MAX_CONCURRENT_UPLOADS).contains(&self.max_concurrent_uploads) { bail!("max concurrent uploads must be between 1 and {MAX_CONCURRENT_UPLOADS}"); } + if self.oss_parallel == 0 { + bail!("oss parallel must be at least 1"); + } + if self.oss_part_size.trim().is_empty() { + bail!("oss part size must be non-empty"); + } Ok(()) } } @@ -1820,11 +1829,8 @@ fn prepare_artifacts(source: &Path, config: &UploadConfig) -> Result<(Artifacts, let data = append_name(source, ".zst")?; let (temporary_data, temporary_file) = exclusive_sibling(&data, ".tmp")?; let output = temporary_file.try_clone()?; - let mut command = Command::new("zstd"); - command - .args(["-q", "-T1", "-3", "-c"]) - .arg(source) - .stdout(Stdio::from(output)); + let mut command = zstd_command(source, config); + command.stdout(Stdio::from(output)); match command_status_with_timeout(&mut command, config.zstd_timeout) { Ok(status) if status.success() => {} Ok(status) => { @@ -1880,6 +1886,15 @@ fn prepare_artifacts(source: &Path, config: &UploadConfig) -> Result<(Artifacts, )) } +fn zstd_command(source: &Path, config: &UploadConfig) -> Command { + let mut command = Command::new("zstd"); + let threads = format!("-T{}", config.zstd_threads); + command + .args(["-q", threads.as_str(), "-3", "-c"]) + .arg(source); + command +} + fn oss_copy_command(source: &str, destination: &str, config: &UploadConfig) -> Command { let mut command = Command::new("aliyun"); command.args([ @@ -1895,6 +1910,11 @@ fn oss_copy_command(source: &str, destination: &str, config: &UploadConfig) -> C &config.region, ]); command + .arg("--parallel") + .arg(config.oss_parallel.to_string()) + .arg("--part-size") + .arg(&config.oss_part_size); + command } fn oss_upload_command(source: &str, destination: &str, config: &UploadConfig) -> Command { @@ -2686,6 +2706,9 @@ mod tests { zstd_timeout: Duration::from_secs(30), oss_timeout: Duration::from_secs(30), max_concurrent_uploads: 2, + zstd_threads: 0, + oss_parallel: 8, + oss_part_size: "32Mi".to_owned(), } } @@ -2707,6 +2730,92 @@ mod tests { .contains("max concurrent uploads must be between 1 and 4")); } + fn command_args(command: &Command) -> Vec { + command + .get_args() + .map(|arg| arg.to_string_lossy().into_owned()) + .collect() + } + + #[test] + fn zstd_command_defaults_to_auto_threads() { + let root = TestDir::new(); + let config = config(root.path()); + let args = command_args(&zstd_command(Path::new("tape.ndjson"), &config)); + assert_eq!(args, ["-q", "-T0", "-3", "-c", "tape.ndjson"]); + } + + #[test] + fn zstd_command_threads_are_configurable() { + let root = TestDir::new(); + let mut config = config(root.path()); + config.zstd_threads = 16; + let args = command_args(&zstd_command(Path::new("tape.ndjson"), &config)); + assert_eq!(args, ["-q", "-T16", "-3", "-c", "tape.ndjson"]); + } + + #[test] + fn oss_copy_command_includes_multipart_tuning() { + let root = TestDir::new(); + let config = config(root.path()); + let args = command_args(&oss_copy_command("src", "dst", &config)); + assert_eq!(&args[..4], ["ossutil", "cp", "src", "dst"]); + assert!(args + .windows(2) + .any(|pair| pair == ["--parallel", "8"])); + assert!(args + .windows(2) + .any(|pair| pair == ["--part-size", "32Mi"])); + } + + #[test] + fn oss_copy_command_tuning_is_configurable() { + let root = TestDir::new(); + let mut config = config(root.path()); + config.oss_parallel = 12; + config.oss_part_size = "64Mi".to_owned(); + let args = command_args(&oss_copy_command("src", "dst", &config)); + assert!(args + .windows(2) + .any(|pair| pair == ["--parallel", "12"])); + assert!(args + .windows(2) + .any(|pair| pair == ["--part-size", "64Mi"])); + } + + #[test] + fn oss_upload_command_keeps_no_clobber_with_tuning() { + let root = TestDir::new(); + let config = config(root.path()); + let args = command_args(&oss_upload_command("src", "dst", &config)); + assert!(args.iter().any(|arg| arg == "--ignore-existing")); + assert!(args + .windows(2) + .any(|pair| pair == ["--parallel", "8"])); + assert!(args + .windows(2) + .any(|pair| pair == ["--part-size", "32Mi"])); + } + + #[test] + fn rejects_invalid_oss_copy_tuning() { + let root = TestDir::new(); + let mut config = config(root.path()); + config.oss_parallel = 0; + assert!(config + .validate() + .unwrap_err() + .to_string() + .contains("oss parallel must be at least 1")); + config.oss_parallel = 8; + config.oss_part_size = " ".to_owned(); + assert!(config + .validate() + .unwrap_err() + .to_string() + .contains("oss part size must be non-empty")); + } + #[test] fn staging_root_creation_is_idempotent_under_concurrency() { let root = TestDir::new();