rbt/materializer/
incremental.rs1use anyhow::{bail, Context, Result};
17use serde::{Deserialize, Serialize};
18use std::fs::{self, File};
19use std::io::Write;
20use std::path::{Path, PathBuf};
21use std::time::{SystemTime, UNIX_EPOCH};
22
23use crate::materializer::stream::{atomic_publish, MaterializeWriteOptions, StreamWriteStats};
24use crate::testing::Assertion;
25use datafusion::physical_plan::SendableRecordBatchStream;
26
27#[derive(Debug, Clone, Serialize, Deserialize)]
29pub struct IncrementalManifest {
30 pub strategy: String,
31 pub parts: Vec<String>,
32 pub total_rows: u64,
33 pub updated_at_ms: u64,
34}
35
36pub fn parts_dir_for_parquet(dest_parquet: &Path) -> PathBuf {
38 let stem = dest_parquet.with_extension("");
39 PathBuf::from(format!("{}.parts", stem.display()))
40}
41
42pub fn manifest_path(parts_dir: &Path) -> PathBuf {
43 parts_dir.join("_rbt_manifest.json")
44}
45
46pub fn load_manifest(parts_dir: &Path) -> Result<IncrementalManifest> {
47 let p = manifest_path(parts_dir);
48 if !p.exists() {
49 return Ok(IncrementalManifest {
50 strategy: "incremental_append".into(),
51 parts: Vec::new(),
52 total_rows: 0,
53 updated_at_ms: 0,
54 });
55 }
56 let s = fs::read_to_string(&p)
57 .with_context(|| format!("E_RBT_INCREMENTAL: read manifest {}", p.display()))?;
58 serde_json::from_str(&s)
59 .with_context(|| format!("E_RBT_INCREMENTAL: parse manifest {}", p.display()))
60}
61
62fn now_ms() -> u64 {
63 SystemTime::now()
64 .duration_since(UNIX_EPOCH)
65 .map(|d| d.as_millis() as u64)
66 .unwrap_or(0)
67}
68
69fn next_part_name(manifest: &IncrementalManifest) -> String {
70 let n = manifest.parts.len() as u64 + 1;
71 format!("part-{n:013}.parquet")
72}
73
74pub async fn materialize_incremental_append_stream(
76 stream: SendableRecordBatchStream,
77 dest_parquet: &Path,
78 opts: &MaterializeWriteOptions,
79 assertions: &[Assertion],
80) -> Result<StreamWriteStats> {
81 let parts_dir = parts_dir_for_parquet(dest_parquet);
82 fs::create_dir_all(&parts_dir).with_context(|| {
83 format!(
84 "E_RBT_INCREMENTAL: mkdir parts {}",
85 parts_dir.display()
86 )
87 })?;
88 let mut manifest = load_manifest(&parts_dir)?;
89 let part_name = next_part_name(&manifest);
90 let part_path = parts_dir.join(&part_name);
91
92 let mut stream = stream;
93 let stats = crate::materializer::stream::write_parquet_stream(
94 &mut stream,
95 &part_path,
96 opts,
97 assertions,
98 )
99 .await
100 .with_context(|| {
101 format!(
102 "E_RBT_INCREMENTAL: write part {}",
103 part_path.display()
104 )
105 })?;
106
107 if stats.rows == 0 {
108 let _ = fs::remove_file(&part_path);
110 return Ok(StreamWriteStats {
111 rows: 0,
112 batches: stats.batches,
113 path: parts_dir,
114 bytes_written: 0,
115 validation: stats.validation,
116 });
117 }
118
119 manifest.parts.push(part_name);
120 manifest.total_rows = manifest.total_rows.saturating_add(stats.rows as u64);
121 manifest.updated_at_ms = now_ms();
122 manifest.strategy = "incremental_append".into();
123 write_manifest(&parts_dir, &manifest)?;
124
125 write_parts_pointer(dest_parquet, &parts_dir)?;
128
129 Ok(StreamWriteStats {
130 rows: stats.rows,
131 batches: stats.batches,
132 path: parts_dir,
133 bytes_written: stats.bytes_written,
134 validation: stats.validation,
135 })
136}
137
138fn write_manifest(parts_dir: &Path, manifest: &IncrementalManifest) -> Result<()> {
139 let p = manifest_path(parts_dir);
140 let partial = p.with_extension("json.partial");
141 {
142 let mut f = File::create(&partial)
143 .with_context(|| format!("E_RBT_INCREMENTAL: create {}", partial.display()))?;
144 writeln!(f, "{}", serde_json::to_string_pretty(manifest)?)?;
145 }
146 atomic_publish(&partial, &p)?;
147 Ok(())
148}
149
150fn write_parts_pointer(dest_parquet: &Path, parts_dir: &Path) -> Result<()> {
152 let pointer = dest_parquet.with_extension("rbt_incremental.json");
153 let body = serde_json::json!({
154 "strategy": "incremental_append",
155 "parts_dir": parts_dir.file_name().and_then(|s| s.to_str()).unwrap_or("parts"),
156 "note": "ref() registers the .parts directory; single-file dest is not rewritten"
157 });
158 fs::write(&pointer, serde_json::to_vec_pretty(&body)?)
159 .with_context(|| format!("E_RBT_INCREMENTAL: write pointer {}", pointer.display()))?;
160 Ok(())
161}
162
163pub fn incremental_ref_path(dest_parquet: &Path) -> PathBuf {
165 let parts = parts_dir_for_parquet(dest_parquet);
166 if parts.is_dir() {
167 parts
168 } else {
169 dest_parquet.to_path_buf()
170 }
171}
172
173pub fn clear_incremental_parts(dest_parquet: &Path) -> Result<()> {
175 let parts = parts_dir_for_parquet(dest_parquet);
176 if parts.exists() {
177 fs::remove_dir_all(&parts).with_context(|| {
178 format!(
179 "E_RBT_INCREMENTAL: clear parts {}",
180 parts.display()
181 )
182 })?;
183 }
184 let pointer = dest_parquet.with_extension("rbt_incremental.json");
185 if pointer.exists() {
186 let _ = fs::remove_file(pointer);
187 }
188 Ok(())
189}
190
191pub fn parse_incremental_strategy(s: &str) -> Result<&'static str> {
193 match s.trim().to_ascii_lowercase().as_str() {
194 "incremental_append" | "append" | "incremental" => Ok("incremental_append"),
195 "table" | "full_refresh" | "full-refresh" => Ok("table"),
196 other => bail!(
197 "E_RBT_INCREMENTAL: unknown materialization '{other}' \
198 (supported: table, incremental_append)"
199 ),
200 }
201}
202
203#[cfg(test)]
204mod tests {
205 use super::*;
206 use datafusion::prelude::SessionContext;
207
208 #[tokio::test]
209 async fn incremental_appends_two_parts() -> Result<()> {
210 let temp = tempfile::tempdir()?;
211 let dest = temp.path().join("stg_x.parquet");
212 let opts = MaterializeWriteOptions::default();
213 let ctx = SessionContext::new();
214
215 let df1 = ctx.sql("SELECT 1 AS id UNION ALL SELECT 2").await?;
216 let s1 = df1.execute_stream().await?;
217 let st1 =
218 materialize_incremental_append_stream(s1, &dest, &opts, &[]).await?;
219 assert_eq!(st1.rows, 2);
220
221 let df2 = ctx.sql("SELECT 3 AS id").await?;
222 let s2 = df2.execute_stream().await?;
223 let st2 =
224 materialize_incremental_append_stream(s2, &dest, &opts, &[]).await?;
225 assert_eq!(st2.rows, 1);
226
227 let parts = parts_dir_for_parquet(&dest);
228 let m = load_manifest(&parts)?;
229 assert_eq!(m.parts.len(), 2);
230 assert_eq!(m.total_rows, 3);
231 assert!(parts.join(&m.parts[0]).exists());
232 assert!(parts.join(&m.parts[1]).exists());
233 Ok(())
234 }
235}