Skip to main content

servers/
row_writer.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;
17
18use ahash::{HashMap as AHashMap, HashMapExt};
19use api::helper::encode_json_value;
20use api::v1::column_data_type_extension::TypeExt;
21use api::v1::helper::time_index_column_schema;
22use api::v1::value::ValueData;
23use api::v1::{
24    ColumnDataType, ColumnDataTypeExtension, ColumnOptions, ColumnSchema, JsonTypeExtension, Row,
25    RowInsertRequest, RowInsertRequests, Rows, SemanticType, Value,
26};
27use arrow_schema::extension::{
28    EXTENSION_TYPE_METADATA_KEY, EXTENSION_TYPE_NAME_KEY, ExtensionType,
29};
30use common_grpc::precision::Precision;
31use common_time::Timestamp;
32use common_time::timestamp::TimeUnit;
33use common_time::timestamp::TimeUnit::Nanosecond;
34use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
35use datatypes::json::JsonSettings;
36use datatypes::value::Value as DataValue;
37use serde::Serialize;
38use snafu::{OptionExt, ResultExt, ensure};
39
40use crate::error::{
41    ConvertScalarValueSnafu, IncompatibleSchemaSnafu, InternalSnafu, Result, RowWriterSnafu,
42    TimePrecisionSnafu, TimestampOverflowSnafu, ToJsonSnafu,
43};
44
45/// The intermediate data structure for building the write request.
46/// It constructs the `schema` and `rows` as all input data row
47/// parsing is completed.
48pub struct TableData {
49    schema: Vec<ColumnSchema>,
50    rows: Vec<Row>,
51    column_indexes: AHashMap<String, usize>,
52}
53
54impl TableData {
55    pub fn new(num_columns: usize, num_rows: usize) -> Self {
56        Self {
57            schema: Vec::with_capacity(num_columns),
58            rows: Vec::with_capacity(num_rows),
59            column_indexes: AHashMap::with_capacity(num_columns),
60        }
61    }
62
63    #[inline]
64    pub fn num_columns(&self) -> usize {
65        self.schema.len()
66    }
67
68    #[inline]
69    pub fn num_rows(&self) -> usize {
70        self.rows.len()
71    }
72
73    #[inline]
74    pub fn alloc_one_row(&self) -> Vec<Value> {
75        vec![Value { value_data: None }; self.num_columns()]
76    }
77
78    #[inline]
79    pub fn add_row(&mut self, values: Vec<Value>) {
80        self.rows.push(Row { values })
81    }
82
83    #[inline]
84    pub fn reserve_rows(&mut self, additional: usize) {
85        self.rows.reserve(additional);
86    }
87
88    pub(crate) fn ensure_column(&mut self, column_schema: ColumnSchema) -> Result<usize> {
89        if let Some(index) = self.column_indexes.get(&column_schema.column_name).copied() {
90            check_schema_number(
91                column_schema.datatype,
92                column_schema.semantic_type,
93                &self.schema[index],
94            )?;
95            return Ok(index);
96        }
97
98        let index = self.schema.len();
99        let name = column_schema.column_name.clone();
100        self.schema.push(column_schema);
101        self.column_indexes.insert(name, index);
102        Ok(index)
103    }
104
105    /// Ensures a default column schema without allocating its name when it already exists.
106    pub(crate) fn ensure_column_by_name(
107        &mut self,
108        name: &str,
109        datatype: ColumnDataType,
110        semantic_type: SemanticType,
111    ) -> Result<usize> {
112        if let Some(index) = self.column_indexes.get(name).copied() {
113            check_schema(datatype, semantic_type, &self.schema[index])?;
114            return Ok(index);
115        }
116
117        self.ensure_column(ColumnSchema {
118            column_name: name.to_string(),
119            datatype: datatype as i32,
120            semantic_type: semantic_type as i32,
121            ..Default::default()
122        })
123    }
124
125    #[allow(dead_code)]
126    pub fn columns(&self) -> &Vec<ColumnSchema> {
127        &self.schema
128    }
129
130    pub fn into_schema_and_rows(self) -> (Vec<ColumnSchema>, Vec<Row>) {
131        (self.schema, self.rows)
132    }
133
134    /// Writes a field value without enforcing that later writes use the same datatype
135    /// as the first-seen schema entry.
136    ///
137    /// The OTLP trace v1 path uses this to preserve raw mixed values inside one request
138    /// so the frontend can reconcile them later against both the full batch and the
139    /// existing table schema.
140    pub fn write_field_unchecked(
141        &mut self,
142        name: impl ToString,
143        datatype: ColumnDataType,
144        value: Option<ValueData>,
145        one_row: &mut Vec<Value>,
146    ) {
147        self.write_column_unchecked(
148            ColumnSchema {
149                column_name: name.to_string(),
150                datatype: datatype as i32,
151                semantic_type: SemanticType::Field as i32,
152                ..Default::default()
153            },
154            value,
155            one_row,
156        );
157    }
158
159    pub fn write_column_unchecked(
160        &mut self,
161        column_schema: ColumnSchema,
162        value: Option<ValueData>,
163        one_row: &mut Vec<Value>,
164    ) {
165        if let Some(index) = self.column_indexes.get(&column_schema.column_name).copied() {
166            one_row[index].value_data = value;
167        } else {
168            let index = self.schema.len();
169            let name = column_schema.column_name.clone();
170            self.schema.push(column_schema);
171            self.column_indexes.insert(name, index);
172            one_row.push(Value { value_data: value });
173        }
174    }
175}
176
177pub struct MultiTableData {
178    table_data_map: HashMap<String, TableData>,
179}
180
181impl Default for MultiTableData {
182    fn default() -> Self {
183        Self::new()
184    }
185}
186
187impl MultiTableData {
188    pub fn new() -> Self {
189        Self {
190            table_data_map: HashMap::new(),
191        }
192    }
193
194    pub fn get_or_default_table_data(
195        &mut self,
196        table_name: impl ToString,
197        num_columns: usize,
198        num_rows: usize,
199    ) -> &mut TableData {
200        self.table_data_map
201            .entry(table_name.to_string())
202            .or_insert_with(|| TableData::new(num_columns, num_rows))
203    }
204
205    pub fn add_table_data(&mut self, table_name: impl ToString, table_data: TableData) {
206        self.table_data_map
207            .insert(table_name.to_string(), table_data);
208    }
209
210    #[allow(dead_code)]
211    pub fn num_tables(&self) -> usize {
212        self.table_data_map.len()
213    }
214
215    /// Returns the request and number of rows in it.
216    pub fn into_row_insert_requests(self) -> (RowInsertRequests, usize) {
217        let mut total_rows = 0;
218        let inserts = self
219            .table_data_map
220            .into_iter()
221            .map(|(table_name, table_data)| {
222                total_rows += table_data.num_rows();
223                let num_columns = table_data.num_columns();
224                let (schema, mut rows) = table_data.into_schema_and_rows();
225                for row in &mut rows {
226                    if num_columns > row.values.len() {
227                        row.values.resize(num_columns, Value { value_data: None });
228                    }
229                }
230
231                RowInsertRequest {
232                    table_name,
233                    rows: Some(Rows { schema, rows }),
234                }
235            })
236            .collect::<Vec<_>>();
237        let row_insert_requests = RowInsertRequests { inserts };
238
239        (row_insert_requests, total_rows)
240    }
241}
242
243/// Write data as tags into the table data.
244pub fn write_tags<K>(
245    table_data: &mut TableData,
246    tags: impl Iterator<Item = (K, String)>,
247    one_row: &mut Vec<Value>,
248) -> Result<()>
249where
250    K: AsRef<str> + Into<String>,
251{
252    let ktv_iter = tags.map(|(k, v)| (k, ColumnDataType::String, Some(ValueData::StringValue(v))));
253    write_by_semantic_type(table_data, SemanticType::Tag, ktv_iter, one_row)
254}
255
256/// Write data as fields into the table data.
257pub fn write_fields(
258    table_data: &mut TableData,
259    fields: impl Iterator<Item = (String, ColumnDataType, Option<ValueData>)>,
260    one_row: &mut Vec<Value>,
261) -> Result<()> {
262    write_by_semantic_type(table_data, SemanticType::Field, fields, one_row)
263}
264
265/// Write data as a tag into the table data.
266pub fn write_tag(
267    table_data: &mut TableData,
268    name: impl ToString,
269    value: impl ToString,
270    one_row: &mut Vec<Value>,
271) -> Result<()> {
272    write_by_semantic_type(
273        table_data,
274        SemanticType::Tag,
275        std::iter::once((
276            name.to_string(),
277            ColumnDataType::String,
278            Some(ValueData::StringValue(value.to_string())),
279        )),
280        one_row,
281    )
282}
283
284/// Write float64 data as a field into the table data.
285pub fn write_f64(
286    table_data: &mut TableData,
287    name: impl ToString,
288    value: f64,
289    one_row: &mut Vec<Value>,
290) -> Result<()> {
291    write_fields(
292        table_data,
293        std::iter::once((
294            name.to_string(),
295            ColumnDataType::Float64,
296            Some(ValueData::F64Value(value)),
297        )),
298        one_row,
299    )
300}
301
302pub(crate) fn build_json_column_schema(name: impl ToString) -> ColumnSchema {
303    ColumnSchema {
304        column_name: name.to_string(),
305        datatype: ColumnDataType::Binary as i32,
306        semantic_type: SemanticType::Field as i32,
307        datatype_extension: Some(ColumnDataTypeExtension {
308            type_ext: Some(TypeExt::JsonType(JsonTypeExtension::JsonBinary.into())),
309        }),
310        ..Default::default()
311    }
312}
313
314pub fn write_json(
315    table_data: &mut TableData,
316    name: impl ToString,
317    value: jsonb::Value,
318    one_row: &mut Vec<Value>,
319) -> Result<()> {
320    write_by_schema(
321        table_data,
322        std::iter::once((
323            build_json_column_schema(name),
324            Some(ValueData::BinaryValue(value.to_vec())),
325        )),
326        one_row,
327    )
328}
329
330pub(crate) fn build_json2_column_schema(name: impl ToString) -> ColumnSchema {
331    let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new(JsonSettings::new_v2())));
332    let mut options = ColumnOptions::default();
333    options.options.insert(
334        EXTENSION_TYPE_NAME_KEY.to_string(),
335        Json2ExtensionType::NAME.to_string(),
336    );
337    if let Some(metadata) = extension.serialize_metadata() {
338        options
339            .options
340            .insert(EXTENSION_TYPE_METADATA_KEY.to_string(), metadata);
341    }
342
343    ColumnSchema {
344        column_name: name.to_string(),
345        datatype: ColumnDataType::Json as i32,
346        semantic_type: SemanticType::Field as i32,
347        options: Some(options),
348        ..Default::default()
349    }
350}
351
352/// Encodes a JSON2 field without constructing its column schema.
353pub(crate) fn encode_json2(value: impl Serialize) -> Result<ValueData> {
354    let json = serde_json::to_value(value).context(ToJsonSnafu)?;
355    let value = JsonSettings::new_v2()
356        .encode(json)
357        .context(ConvertScalarValueSnafu)?;
358    let DataValue::Json(value) = value else {
359        return InternalSnafu {
360            err_msg: "JSON2 encoding returned a non-JSON value",
361        }
362        .fail();
363    };
364    Ok(ValueData::JsonValue(encode_json_value(*value)))
365}
366
367pub(crate) fn write_by_schema(
368    table_data: &mut TableData,
369    kv_iter: impl Iterator<Item = (ColumnSchema, Option<ValueData>)>,
370    one_row: &mut Vec<Value>,
371) -> Result<()> {
372    let TableData {
373        schema,
374        column_indexes,
375        ..
376    } = table_data;
377
378    for (column_schema, value) in kv_iter {
379        let index = column_indexes.get(&column_schema.column_name);
380        if let Some(index) = index {
381            check_schema_number(
382                column_schema.datatype,
383                column_schema.semantic_type,
384                &schema[*index],
385            )?;
386            one_row[*index].value_data = value;
387        } else {
388            let index = schema.len();
389            let key = column_schema.column_name.clone();
390            schema.push(column_schema);
391            column_indexes.insert(key, index);
392            one_row.push(Value { value_data: value });
393        }
394    }
395
396    Ok(())
397}
398
399fn write_by_semantic_type<K>(
400    table_data: &mut TableData,
401    semantic_type: SemanticType,
402    ktv_iter: impl Iterator<Item = (K, ColumnDataType, Option<ValueData>)>,
403    one_row: &mut Vec<Value>,
404) -> Result<()>
405where
406    K: AsRef<str> + Into<String>,
407{
408    let TableData {
409        schema,
410        column_indexes,
411        ..
412    } = table_data;
413
414    for (name, datatype, value) in ktv_iter {
415        let index = column_indexes.get(name.as_ref()).copied();
416        if let Some(index) = index {
417            check_schema(datatype, semantic_type, &schema[index])?;
418            one_row[index].value_data = value;
419        } else {
420            let index = schema.len();
421            let name = name.into();
422            schema.push(ColumnSchema {
423                column_name: name.clone(),
424                datatype: datatype as i32,
425                semantic_type: semantic_type as i32,
426                ..Default::default()
427            });
428            column_indexes.insert(name, index);
429            one_row.push(Value { value_data: value });
430        }
431    }
432
433    Ok(())
434}
435
436/// Write timestamp data as milliseconds into the table data.
437pub fn write_ts_to_millis(
438    table_data: &mut TableData,
439    name: impl ToString,
440    ts: Option<i64>,
441    precision: Precision,
442    one_row: &mut Vec<Value>,
443) -> Result<()> {
444    write_ts_to(
445        table_data,
446        name,
447        ts,
448        precision,
449        TimestampType::Millis,
450        one_row,
451    )
452}
453
454/// Write timestamp data as nanoseconds into the table data.
455pub fn write_ts_to_nanos(
456    table_data: &mut TableData,
457    name: impl ToString,
458    ts: Option<i64>,
459    precision: Precision,
460    one_row: &mut Vec<Value>,
461) -> Result<()> {
462    write_ts_to(
463        table_data,
464        name,
465        ts,
466        precision,
467        TimestampType::Nanos,
468        one_row,
469    )
470}
471
472enum TimestampType {
473    Millis,
474    Nanos,
475}
476
477fn write_ts_to(
478    table_data: &mut TableData,
479    name: impl ToString,
480    ts: Option<i64>,
481    precision: Precision,
482    ts_type: TimestampType,
483    one_row: &mut Vec<Value>,
484) -> Result<()> {
485    let TableData {
486        schema,
487        column_indexes,
488        ..
489    } = table_data;
490    let name = name.to_string();
491
492    let ts = match ts {
493        Some(timestamp) => match ts_type {
494            TimestampType::Millis => precision.to_millis(timestamp),
495            TimestampType::Nanos => precision.to_nanos(timestamp),
496        }
497        .with_context(|| TimestampOverflowSnafu {
498            error: format!(
499                "timestamp {} overflow with precision {}",
500                timestamp, precision
501            ),
502        })?,
503        None => {
504            let timestamp = Timestamp::current_time(Nanosecond);
505            let unit: TimeUnit = precision.try_into().context(RowWriterSnafu)?;
506            let timestamp = timestamp
507                .convert_to(unit)
508                .with_context(|| TimePrecisionSnafu {
509                    name: precision.to_string(),
510                })?
511                .into();
512            match ts_type {
513                TimestampType::Millis => precision.to_millis(timestamp),
514                TimestampType::Nanos => precision.to_nanos(timestamp),
515            }
516            .with_context(|| TimestampOverflowSnafu {
517                error: format!(
518                    "timestamp {} overflow with precision {}",
519                    timestamp, precision
520                ),
521            })?
522        }
523    };
524
525    let (datatype, ts) = match ts_type {
526        TimestampType::Millis => (
527            ColumnDataType::TimestampMillisecond,
528            ValueData::TimestampMillisecondValue(ts),
529        ),
530        TimestampType::Nanos => (
531            ColumnDataType::TimestampNanosecond,
532            ValueData::TimestampNanosecondValue(ts),
533        ),
534    };
535
536    let index = column_indexes.get(&name);
537    if let Some(index) = index {
538        check_schema(datatype, SemanticType::Timestamp, &schema[*index])?;
539        one_row[*index].value_data = Some(ts);
540    } else {
541        let index = schema.len();
542        schema.push(time_index_column_schema(&name, datatype));
543        column_indexes.insert(name, index);
544        one_row.push(ts.into())
545    }
546
547    Ok(())
548}
549
550fn check_schema(
551    datatype: ColumnDataType,
552    semantic_type: SemanticType,
553    schema: &ColumnSchema,
554) -> Result<()> {
555    check_schema_number(datatype as i32, semantic_type as i32, schema)
556}
557
558fn check_schema_number(datatype: i32, semantic_type: i32, schema: &ColumnSchema) -> Result<()> {
559    ensure!(
560        schema.datatype == datatype,
561        IncompatibleSchemaSnafu {
562            column_name: &schema.column_name,
563            datatype: "datatype",
564            expected: schema.datatype,
565            actual: datatype,
566        }
567    );
568
569    ensure!(
570        schema.semantic_type == semantic_type,
571        IncompatibleSchemaSnafu {
572            column_name: &schema.column_name,
573            datatype: "semantic_type",
574            expected: schema.semantic_type,
575            actual: semantic_type,
576        }
577    );
578
579    Ok(())
580}