Skip to main content

iceberg/spec/
view_version.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18/*!
19 * View Versions!
20 */
21use std::collections::HashMap;
22use std::sync::Arc;
23
24use _serde::ViewVersionV1;
25use chrono::{DateTime, Utc};
26use serde::{Deserialize, Serialize};
27use typed_builder::TypedBuilder;
28
29use super::INITIAL_VIEW_VERSION_ID;
30use super::view_metadata::ViewVersionLog;
31use crate::catalog::NamespaceIdent;
32use crate::error::{Result, invalid_data, timestamp_ms_to_utc};
33use crate::spec::{SchemaId, SchemaRef, ViewMetadata};
34
35/// Reference to [`ViewVersion`].
36pub type ViewVersionRef = Arc<ViewVersion>;
37
38/// Alias for the integer type used for view version ids.
39pub type ViewVersionId = i32;
40
41#[derive(Debug, PartialEq, Eq, Clone, Serialize, Deserialize, TypedBuilder)]
42#[serde(from = "ViewVersionV1", into = "ViewVersionV1")]
43#[builder(field_defaults(setter(prefix = "with_")))]
44/// A view versions represents the definition of a view at a specific point in time.
45pub struct ViewVersion {
46    /// A unique long ID
47    #[builder(default = INITIAL_VIEW_VERSION_ID)]
48    version_id: ViewVersionId,
49    /// ID of the schema for the view version
50    schema_id: SchemaId,
51    /// Timestamp when the version was created (ms from epoch)
52    timestamp_ms: i64,
53    /// A string to string map of summary metadata about the version
54    #[builder(default = HashMap::new())]
55    summary: HashMap<String, String>,
56    /// A list of representations for the view definition.
57    representations: ViewRepresentations,
58    /// Catalog name to use when a reference in the SELECT does not contain a catalog
59    #[builder(default = None)]
60    default_catalog: Option<String>,
61    /// Namespace to use when a reference in the SELECT is a single identifier
62    default_namespace: NamespaceIdent,
63}
64
65impl ViewVersion {
66    /// Get the version id of this view version.
67    #[inline]
68    pub fn version_id(&self) -> ViewVersionId {
69        self.version_id
70    }
71
72    /// Get the schema id of this view version.
73    #[inline]
74    pub fn schema_id(&self) -> SchemaId {
75        self.schema_id
76    }
77
78    /// Get the timestamp of when the view version was created
79    #[inline]
80    pub fn timestamp(&self) -> Result<DateTime<Utc>> {
81        timestamp_ms_to_utc(self.timestamp_ms)
82    }
83
84    /// Get the timestamp of when the view version was created in milliseconds since epoch
85    #[inline]
86    pub fn timestamp_ms(&self) -> i64 {
87        self.timestamp_ms
88    }
89
90    /// Get summary of the view version
91    #[inline]
92    pub fn summary(&self) -> &HashMap<String, String> {
93        &self.summary
94    }
95
96    /// Get this views representations
97    #[inline]
98    pub fn representations(&self) -> &ViewRepresentations {
99        &self.representations
100    }
101
102    /// Get the default catalog for this view version
103    #[inline]
104    pub fn default_catalog(&self) -> Option<&String> {
105        self.default_catalog.as_ref()
106    }
107
108    /// Get the default namespace to use when a reference in the SELECT is a single identifier
109    #[inline]
110    pub fn default_namespace(&self) -> &NamespaceIdent {
111        &self.default_namespace
112    }
113
114    /// Get the schema of this snapshot.
115    pub fn schema(&self, view_metadata: &ViewMetadata) -> Result<SchemaRef> {
116        view_metadata
117            .schema_by_id(self.schema_id())
118            .ok_or_else(|| invalid_data!("Schema with id {} not found", self.schema_id()))
119            .cloned()
120    }
121
122    /// Retrieve the history log entry for this view version.
123    pub(crate) fn log(&self) -> ViewVersionLog {
124        ViewVersionLog::new(self.version_id, self.timestamp_ms)
125    }
126
127    /// Check if this view version behaves the same as another view spec.
128    ///
129    /// Returns true if the view version is equal to the other view version
130    /// with `timestamp_ms` and `version_id` ignored. The following must be identical:
131    /// * Summary (all of them)
132    /// * Representations
133    /// * Default Catalog
134    /// * Default Namespace
135    /// * The Schema ID
136    pub(crate) fn behaves_identical_to(&self, other: &Self) -> bool {
137        self.summary == other.summary
138            && self.representations == other.representations
139            && self.default_catalog == other.default_catalog
140            && self.default_namespace == other.default_namespace
141            && self.schema_id == other.schema_id
142    }
143
144    /// Change the version id of this view version.
145    pub fn with_version_id(self, version_id: i32) -> Self {
146        Self { version_id, ..self }
147    }
148
149    /// Change the schema id of this view version.
150    pub fn with_schema_id(self, schema_id: SchemaId) -> Self {
151        Self { schema_id, ..self }
152    }
153}
154
155/// A list of view representations.
156#[derive(Debug, PartialEq, Eq, Clone, Serialize, Deserialize)]
157pub struct ViewRepresentations(pub(crate) Vec<ViewRepresentation>);
158
159impl ViewRepresentations {
160    #[inline]
161    /// Get the number of representations
162    pub fn len(&self) -> usize {
163        self.0.len()
164    }
165
166    #[inline]
167    /// Check if there are no representations
168    pub fn is_empty(&self) -> bool {
169        self.0.is_empty()
170    }
171
172    /// Get an iterator over the representations
173    pub fn iter(&self) -> impl ExactSizeIterator<Item = &'_ ViewRepresentation> {
174        self.0.iter()
175    }
176}
177
178// Iterator for ViewRepresentations
179impl IntoIterator for ViewRepresentations {
180    type Item = ViewRepresentation;
181    type IntoIter = std::vec::IntoIter<Self::Item>;
182
183    fn into_iter(self) -> Self::IntoIter {
184        self.0.into_iter()
185    }
186}
187
188#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
189#[serde(tag = "type")]
190/// View definitions can be represented in multiple ways.
191/// Representations are documented ways to express a view definition.
192// ToDo: Make unique per Dialect
193pub enum ViewRepresentation {
194    #[serde(rename = "sql")]
195    /// The SQL representation stores the view definition as a SQL SELECT,
196    Sql(SqlViewRepresentation),
197}
198
199#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
200#[serde(rename_all = "kebab-case")]
201/// The SQL representation stores the view definition as a SQL SELECT,
202/// with metadata such as the SQL dialect.
203pub struct SqlViewRepresentation {
204    #[serde(rename = "sql")]
205    /// The SQL SELECT statement that defines the view.
206    pub sql: String,
207    #[serde(rename = "dialect")]
208    /// The dialect of the sql SELECT statement (e.g., "trino" or "spark")
209    pub dialect: String,
210}
211
212pub(super) mod _serde {
213    /// This is a helper module that defines types to help with serialization/deserialization.
214    /// For deserialization the input first gets read into the [`ViewVersionV1`] struct.
215    /// and then converted into the [Snapshot] struct. Serialization works the other way around.
216    /// [ViewVersionV1] are internal struct that are only used for serialization and deserialization.
217    use serde::{Deserialize, Serialize};
218
219    use super::{ViewRepresentation, ViewRepresentations, ViewVersion};
220    use crate::catalog::NamespaceIdent;
221
222    #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
223    #[serde(rename_all = "kebab-case")]
224    /// Defines the structure of a v1 view version for serialization/deserialization
225    pub(crate) struct ViewVersionV1 {
226        pub version_id: i32,
227        pub schema_id: i32,
228        pub timestamp_ms: i64,
229        pub summary: std::collections::HashMap<String, String>,
230        pub representations: Vec<ViewRepresentation>,
231        #[serde(skip_serializing_if = "Option::is_none")]
232        pub default_catalog: Option<String>,
233        pub default_namespace: NamespaceIdent,
234    }
235
236    impl From<ViewVersionV1> for ViewVersion {
237        fn from(v1: ViewVersionV1) -> Self {
238            ViewVersion {
239                version_id: v1.version_id,
240                schema_id: v1.schema_id,
241                timestamp_ms: v1.timestamp_ms,
242                summary: v1.summary,
243                representations: ViewRepresentations(v1.representations),
244                default_catalog: v1.default_catalog,
245                default_namespace: v1.default_namespace,
246            }
247        }
248    }
249
250    impl From<ViewVersion> for ViewVersionV1 {
251        fn from(v1: ViewVersion) -> Self {
252            ViewVersionV1 {
253                version_id: v1.version_id,
254                schema_id: v1.schema_id,
255                timestamp_ms: v1.timestamp_ms,
256                summary: v1.summary,
257                representations: v1.representations.0,
258                default_catalog: v1.default_catalog,
259                default_namespace: v1.default_namespace,
260            }
261        }
262    }
263}
264
265impl From<SqlViewRepresentation> for ViewRepresentation {
266    fn from(sql: SqlViewRepresentation) -> Self {
267        ViewRepresentation::Sql(sql)
268    }
269}
270
271#[cfg(test)]
272mod tests {
273    use chrono::{TimeZone, Utc};
274
275    use crate::NamespaceIdent;
276    use crate::spec::ViewRepresentations;
277    use crate::spec::view_version::_serde::ViewVersionV1;
278    use crate::spec::view_version::ViewVersion;
279
280    #[test]
281    fn view_version() {
282        let record = serde_json::json!(
283        {
284            "version-id" : 1,
285            "timestamp-ms" : 1573518431292i64,
286            "schema-id" : 1,
287            "default-catalog" : "prod",
288            "default-namespace" : [ "default" ],
289            "summary" : {
290              "engine-name" : "Spark",
291              "engineVersion" : "3.3.2"
292            },
293            "representations" : [ {
294              "type" : "sql",
295              "sql" : "SELECT\n    COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2",
296              "dialect" : "spark"
297            } ]
298          }
299        );
300
301        let result: ViewVersion = serde_json::from_value::<ViewVersionV1>(record.clone())
302            .unwrap()
303            .into();
304
305        // Roundtrip
306        assert_eq!(serde_json::to_value(result.clone()).unwrap(), record);
307
308        assert_eq!(result.version_id(), 1);
309        assert_eq!(
310            result.timestamp().unwrap(),
311            Utc.timestamp_millis_opt(1573518431292).unwrap()
312        );
313        assert_eq!(result.schema_id(), 1);
314        assert_eq!(result.default_catalog, Some("prod".to_string()));
315        assert_eq!(result.summary(), &{
316            let mut map = std::collections::HashMap::new();
317            map.insert("engine-name".to_string(), "Spark".to_string());
318            map.insert("engineVersion".to_string(), "3.3.2".to_string());
319            map
320        });
321        assert_eq!(
322            result.representations().to_owned(),
323            ViewRepresentations(vec![super::ViewRepresentation::Sql(
324                super::SqlViewRepresentation {
325                    sql: "SELECT\n    COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
326                        .to_string(),
327                    dialect: "spark".to_string(),
328                },
329            )])
330        );
331        assert_eq!(result.default_namespace.inner(), vec![
332            "default".to_string()
333        ]);
334    }
335
336    #[test]
337    fn test_behaves_identical_to() {
338        let view_version = ViewVersion::builder()
339            .with_version_id(1)
340            .with_schema_id(1)
341            .with_timestamp_ms(1573518431292)
342            .with_summary({
343                let mut map = std::collections::HashMap::new();
344                map.insert("engine-name".to_string(), "Spark".to_string());
345                map.insert("engineVersion".to_string(), "3.3.2".to_string());
346                map
347            })
348            .with_representations(ViewRepresentations(vec![super::ViewRepresentation::Sql(
349                super::SqlViewRepresentation {
350                    sql: "SELECT\n    COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
351                        .to_string(),
352                    dialect: "spark".to_string(),
353                },
354            )]))
355            .with_default_catalog(Some("prod".to_string()))
356            .with_default_namespace(NamespaceIdent::new("default".to_string()))
357            .build();
358
359        let mut identical_view_version = view_version.clone();
360        identical_view_version.version_id = 2;
361        identical_view_version.timestamp_ms = 1573518431293;
362
363        let different_view_version = ViewVersion::builder()
364            .with_version_id(view_version.version_id())
365            .with_schema_id(view_version.schema_id())
366            .with_timestamp_ms(view_version.timestamp_ms())
367            .with_summary(view_version.summary().clone())
368            .with_representations(ViewRepresentations(vec![super::ViewRepresentation::Sql(
369                super::SqlViewRepresentation {
370                    sql: "SELECT * from events".to_string(),
371                    dialect: "spark".to_string(),
372                },
373            )]))
374            .with_default_catalog(view_version.default_catalog().cloned())
375            .with_default_namespace(view_version.default_namespace().clone())
376            .build();
377
378        assert!(view_version.behaves_identical_to(&identical_view_version));
379        assert!(!view_version.behaves_identical_to(&different_view_version));
380    }
381}