Skip to main content

spark_connect/
storage.rs

1//! Storage-level presets, mirroring `pyspark.StorageLevel`.
2//!
3//! `df.persist(...)` requires a valid storage level. The proto `StorageLevel`
4//! default is all-false (equivalent to `NONE`), which the server rejects with
5//! "StorageLevel is null or invalid" — so callers should use one of the named
6//! presets below rather than `StorageLevel::default()`.
7
8use spark_connect_proto::StorageLevel;
9
10/// Named storage-level presets on [`StorageLevel`], matching the constants on
11/// `pyspark.StorageLevel` (`MEMORY_AND_DISK`, `DISK_ONLY`, `OFF_HEAP`, ...).
12///
13/// Bring this trait into scope to call e.g. `StorageLevel::memory_and_disk()`.
14pub trait StorageLevelExt {
15    /// No storage (`NONE`). Not a valid level to `persist(...)` with.
16    fn none() -> StorageLevel;
17    /// Disk only, 1 replica (`DISK_ONLY`).
18    fn disk_only() -> StorageLevel;
19    /// Disk only, 2 replicas (`DISK_ONLY_2`).
20    fn disk_only_2() -> StorageLevel;
21    /// Disk only, 3 replicas (`DISK_ONLY_3`).
22    fn disk_only_3() -> StorageLevel;
23    /// Memory only, 1 replica (`MEMORY_ONLY`).
24    fn memory_only() -> StorageLevel;
25    /// Memory only, 2 replicas (`MEMORY_ONLY_2`).
26    fn memory_only_2() -> StorageLevel;
27    /// Memory and disk, 1 replica (`MEMORY_AND_DISK`).
28    fn memory_and_disk() -> StorageLevel;
29    /// Memory and disk, 2 replicas (`MEMORY_AND_DISK_2`).
30    fn memory_and_disk_2() -> StorageLevel;
31    /// Memory and disk, deserialized, 1 replica (`MEMORY_AND_DISK_DESER`) — the
32    /// default level used by `cache()`.
33    fn memory_and_disk_deser() -> StorageLevel;
34    /// Off-heap (`OFF_HEAP`).
35    fn off_heap() -> StorageLevel;
36}
37
38fn level(
39    use_disk: bool,
40    use_memory: bool,
41    use_off_heap: bool,
42    deserialized: bool,
43    replication: i32,
44) -> StorageLevel {
45    StorageLevel {
46        use_disk,
47        use_memory,
48        use_off_heap,
49        deserialized,
50        replication,
51    }
52}
53
54impl StorageLevelExt for StorageLevel {
55    fn none() -> StorageLevel {
56        level(false, false, false, false, 1)
57    }
58    fn disk_only() -> StorageLevel {
59        level(true, false, false, false, 1)
60    }
61    fn disk_only_2() -> StorageLevel {
62        level(true, false, false, false, 2)
63    }
64    fn disk_only_3() -> StorageLevel {
65        level(true, false, false, false, 3)
66    }
67    fn memory_only() -> StorageLevel {
68        level(false, true, false, false, 1)
69    }
70    fn memory_only_2() -> StorageLevel {
71        level(false, true, false, false, 2)
72    }
73    fn memory_and_disk() -> StorageLevel {
74        level(true, true, false, false, 1)
75    }
76    fn memory_and_disk_2() -> StorageLevel {
77        level(true, true, false, false, 2)
78    }
79    fn memory_and_disk_deser() -> StorageLevel {
80        level(true, true, false, true, 1)
81    }
82    fn off_heap() -> StorageLevel {
83        level(true, true, true, false, 1)
84    }
85}
86
87#[cfg(test)]
88mod tests {
89    use super::*;
90
91    #[test]
92    fn presets_match_reference_values() {
93        let s = StorageLevel::memory_and_disk_deser();
94        assert!(s.use_disk && s.use_memory && s.deserialized && !s.use_off_heap);
95        assert_eq!(s.replication, 1);
96
97        let s = StorageLevel::memory_and_disk();
98        assert!(s.use_disk && s.use_memory && !s.deserialized && !s.use_off_heap);
99        assert_eq!(s.replication, 1);
100
101        let s = StorageLevel::disk_only_3();
102        assert!(s.use_disk && !s.use_memory);
103        assert_eq!(s.replication, 3);
104
105        let s = StorageLevel::off_heap();
106        assert!(s.use_off_heap && s.use_disk && s.use_memory);
107
108        // NONE is all-false (replication 1) — invalid to persist with.
109        let s = StorageLevel::none();
110        assert!(!s.use_disk && !s.use_memory && !s.use_off_heap && !s.deserialized);
111    }
112
113    #[test]
114    fn every_preset_matches_its_flags_and_replication() {
115        // Exercise the remaining presets so all StorageLevelExt methods are covered.
116        let s = StorageLevel::disk_only();
117        assert!(s.use_disk && !s.use_memory && !s.use_off_heap && !s.deserialized);
118        assert_eq!(s.replication, 1);
119
120        let s = StorageLevel::disk_only_2();
121        assert!(s.use_disk && !s.use_memory);
122        assert_eq!(s.replication, 2);
123
124        let s = StorageLevel::memory_only();
125        assert!(s.use_memory && !s.use_disk && !s.use_off_heap);
126        assert_eq!(s.replication, 1);
127
128        let s = StorageLevel::memory_only_2();
129        assert!(s.use_memory && !s.use_disk);
130        assert_eq!(s.replication, 2);
131
132        let s = StorageLevel::memory_and_disk_2();
133        assert!(s.use_disk && s.use_memory);
134        assert_eq!(s.replication, 2);
135    }
136}