Skip to main content

pipeline/manager/
pipeline_operator.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16use std::sync::{Arc, RwLock};
17use std::time::{Duration, Instant};
18
19use api::v1::CreateTableExpr;
20use catalog::{CatalogManagerRef, RegisterSystemTableRequest};
21use common_catalog::consts::{DEFAULT_PRIVATE_SCHEMA_NAME, default_engine};
22use common_frontend::metrics;
23use common_meta::rpc::ddl::TriggerReason;
24use common_telemetry::info;
25use common_time::FOREVER;
26use datatypes::timestamp::TimestampNanosecond;
27use futures::FutureExt;
28use operator::insert::InserterRef;
29use operator::statement::StatementExecutorRef;
30use query::QueryEngineRef;
31use session::context::QueryContextRef;
32use snafu::{OptionExt, ResultExt};
33use table::TableRef;
34use table::requests::TTL_KEY;
35
36use crate::Pipeline;
37use crate::error::{CatalogSnafu, CreateTableSnafu, PipelineTableNotFoundSnafu, Result};
38use crate::manager::{PipelineInfo, PipelineTableRef, PipelineVersion};
39use crate::metrics::{
40    METRIC_PIPELINE_CREATE_HISTOGRAM, METRIC_PIPELINE_DELETE_HISTOGRAM,
41    METRIC_PIPELINE_RETRIEVE_HISTOGRAM,
42};
43use crate::options::PipelineOptions;
44use crate::table::{PIPELINE_TABLE_NAME, PipelineTable};
45
46/// PipelineOperator is responsible for managing pipelines.
47/// It provides the ability to:
48/// - Create a pipeline table if it does not exist
49/// - Get a pipeline from the pipeline table
50/// - Insert a pipeline into the pipeline table
51/// - Compile a pipeline
52/// - Add a pipeline table to the cache
53/// - Get a pipeline table from the cache
54pub struct PipelineOperator {
55    inserter: InserterRef,
56    statement_executor: StatementExecutorRef,
57    catalog_manager: CatalogManagerRef,
58    query_engine: QueryEngineRef,
59    tables: RwLock<HashMap<String, PipelineTableRef>>,
60    cache_ttl: Duration,
61}
62
63impl PipelineOperator {
64    /// Create a table request for the pipeline table.
65    fn create_table_request(&self, catalog: &str) -> RegisterSystemTableRequest {
66        let (time_index, primary_keys, column_defs) = PipelineTable::build_pipeline_schema();
67
68        let mut table_options = HashMap::new();
69        table_options.insert(TTL_KEY.to_string(), FOREVER.to_string());
70
71        let create_table_expr = CreateTableExpr {
72            catalog_name: catalog.to_string(),
73            schema_name: DEFAULT_PRIVATE_SCHEMA_NAME.to_string(),
74            table_name: PIPELINE_TABLE_NAME.to_string(),
75            desc: "GreptimeDB pipeline table for Log".to_string(),
76            column_defs,
77            time_index,
78            primary_keys,
79            create_if_not_exists: true,
80            table_options,
81            table_id: None, // Should and will be assigned by Meta.
82            engine: default_engine().to_string(),
83        };
84
85        RegisterSystemTableRequest {
86            create_table_expr,
87            open_hook: None,
88        }
89    }
90
91    fn add_pipeline_table_to_cache(&self, catalog: &str, table: TableRef) {
92        let mut tables = self.tables.write().unwrap();
93        if tables.contains_key(catalog) {
94            return;
95        }
96        tables.insert(
97            catalog.to_string(),
98            Arc::new(PipelineTable::new(
99                self.inserter.clone(),
100                self.statement_executor.clone(),
101                table,
102                self.query_engine.clone(),
103                self.cache_ttl,
104            )),
105        );
106    }
107
108    async fn create_pipeline_table_if_not_exists(&self, ctx: QueryContextRef) -> Result<()> {
109        let catalog = ctx.current_catalog();
110
111        // exist in cache
112        if self.get_pipeline_table_from_cache(catalog).is_some() {
113            return Ok(());
114        }
115
116        let RegisterSystemTableRequest {
117            create_table_expr: mut expr,
118            open_hook: _,
119        } = self.create_table_request(catalog);
120
121        // exist in catalog, just open
122        if let Some(table) = self
123            .catalog_manager
124            .table(
125                &expr.catalog_name,
126                &expr.schema_name,
127                &expr.table_name,
128                Some(&ctx),
129            )
130            .await
131            .context(CatalogSnafu)?
132        {
133            self.add_pipeline_table_to_cache(catalog, table);
134            return Ok(());
135        }
136
137        // create table
138        self.statement_executor
139            .create_table_inner(&mut expr, None, ctx.clone(), TriggerReason::AutoCreate)
140            .await
141            .context(CreateTableSnafu)?;
142
143        let schema = &expr.schema_name;
144        let table_name = &expr.table_name;
145
146        // get from catalog
147        let table = self
148            .catalog_manager
149            .table(catalog, schema, table_name, Some(&ctx))
150            .await
151            .context(CatalogSnafu)?
152            .context(PipelineTableNotFoundSnafu)?;
153
154        info!(
155            "Created pipelines table {} with table id {}.",
156            table.table_info().full_table_name(),
157            table.table_info().table_id()
158        );
159
160        // put to cache
161        self.add_pipeline_table_to_cache(catalog, table);
162
163        Ok(())
164    }
165
166    /// Get a pipeline table from the cache.
167    pub fn get_pipeline_table_from_cache(&self, catalog: &str) -> Option<PipelineTableRef> {
168        let table = self.tables.read().unwrap().get(catalog).cloned();
169        metrics::record_cache_lookup("pipeline_table", table.is_some());
170        table
171    }
172}
173
174impl PipelineOperator {
175    /// Create a new PipelineOperator.
176    pub fn new(
177        inserter: InserterRef,
178        statement_executor: StatementExecutorRef,
179        catalog_manager: CatalogManagerRef,
180        query_engine: QueryEngineRef,
181        options: &PipelineOptions,
182    ) -> Self {
183        Self {
184            inserter,
185            statement_executor,
186            catalog_manager,
187            tables: RwLock::new(HashMap::new()),
188            query_engine,
189            cache_ttl: options.cache_ttl,
190        }
191    }
192
193    /// Get a pipeline from the pipeline table.
194    pub async fn get_pipeline(
195        &self,
196        query_ctx: QueryContextRef,
197        name: &str,
198        version: PipelineVersion,
199    ) -> Result<Arc<Pipeline>> {
200        let schema = query_ctx.current_schema();
201        self.create_pipeline_table_if_not_exists(query_ctx.clone())
202            .await?;
203
204        let timer = Instant::now();
205        self.get_pipeline_table_from_cache(query_ctx.current_catalog())
206            .context(PipelineTableNotFoundSnafu)?
207            .get_pipeline(&schema, name, version)
208            .inspect(|re| {
209                METRIC_PIPELINE_RETRIEVE_HISTOGRAM
210                    .with_label_values(&[&re.is_ok().to_string()])
211                    .observe(timer.elapsed().as_secs_f64())
212            })
213            .await
214    }
215
216    /// Get a original pipeline by name.
217    pub async fn get_pipeline_str(
218        &self,
219        name: &str,
220        version: PipelineVersion,
221        query_ctx: QueryContextRef,
222    ) -> Result<(String, TimestampNanosecond)> {
223        let schema = query_ctx.current_schema();
224        self.create_pipeline_table_if_not_exists(query_ctx.clone())
225            .await?;
226
227        let timer = Instant::now();
228        self.get_pipeline_table_from_cache(query_ctx.current_catalog())
229            .context(PipelineTableNotFoundSnafu)?
230            .get_pipeline_str(&schema, name, version)
231            .inspect(|re| {
232                METRIC_PIPELINE_RETRIEVE_HISTOGRAM
233                    .with_label_values(&[&re.is_ok().to_string()])
234                    .observe(timer.elapsed().as_secs_f64())
235            })
236            .await
237            .map(|p| (p.content, p.version))
238    }
239
240    /// Insert a pipeline into the pipeline table.
241    pub async fn insert_pipeline(
242        &self,
243        name: &str,
244        content_type: &str,
245        pipeline: &str,
246        query_ctx: QueryContextRef,
247    ) -> Result<PipelineInfo> {
248        self.create_pipeline_table_if_not_exists(query_ctx.clone())
249            .await?;
250
251        let timer = Instant::now();
252        self.get_pipeline_table_from_cache(query_ctx.current_catalog())
253            .context(PipelineTableNotFoundSnafu)?
254            .insert_and_compile(name, content_type, pipeline)
255            .inspect(|re| {
256                METRIC_PIPELINE_CREATE_HISTOGRAM
257                    .with_label_values(&[&re.is_ok().to_string()])
258                    .observe(timer.elapsed().as_secs_f64())
259            })
260            .await
261    }
262
263    /// Delete a pipeline by name from pipeline table.
264    pub async fn delete_pipeline(
265        &self,
266        name: &str,
267        version: PipelineVersion,
268        query_ctx: QueryContextRef,
269    ) -> Result<Option<()>> {
270        // trigger load pipeline table
271        self.create_pipeline_table_if_not_exists(query_ctx.clone())
272            .await?;
273
274        let timer = Instant::now();
275        self.get_pipeline_table_from_cache(query_ctx.current_catalog())
276            .context(PipelineTableNotFoundSnafu)?
277            .delete_pipeline(name, version)
278            .inspect(|re| {
279                METRIC_PIPELINE_DELETE_HISTOGRAM
280                    .with_label_values(&[&re.is_ok().to_string()])
281                    .observe(timer.elapsed().as_secs_f64())
282            })
283            .await
284    }
285
286    /// Compile a pipeline.
287    pub fn build_pipeline(pipeline: &str) -> Result<Pipeline> {
288        PipelineTable::compile_pipeline(pipeline)
289    }
290}