Skip to main content

rbt/materializer/
incremental.rs

1//! Incremental append materialization for parquet models.
2//!
3//! **Honest scope:** append-only **part files** under `{model}.parts/`, not row-level MERGE.
4//! Full-refresh `table` materialization still overwrites a single file.
5//!
6//! Layout:
7//! ```text
8//! lake/silver/stg_events.parts/
9//!   part-0000000000001.parquet
10//!   part-0000000000002.parquet
11//!   _rbt_manifest.json
12//! ```
13//!
14//! Downstream `ref()` registers the **parts directory** as a multi-file parquet table.
15
16use 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/// Manifest describing incremental parts for a model.
28#[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
36/// Parts directory sibling to a flat parquet path: `foo.parquet` → `foo.parts/`.
37pub 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
74/// Stream-write a new part and update the manifest (append-only).
75pub 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        // Empty increment: remove empty part if created, leave manifest unchanged.
109        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    // Optional convenience: also refresh a single-file view for tools that expect dest_parquet.
126    // We do **not** rewrite the full union here (would defeat incremental). Point dest at parts via symlink if supported.
127    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
150/// Write a tiny pointer file next to the logical model path so operators know this is incremental.
151fn 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
163/// Path to register for `ref()` when incremental: the parts directory.
164pub 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
173/// Full-refresh wipe of incremental parts (when model switches back to table, or explicit).
174pub 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
191/// Validate incremental frontmatter hints.
192pub 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}