Skip to main content

pipeline/etl/
value.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::BTreeMap;
16
17use api::v1::ColumnDataType;
18use api::v1::value::ValueData;
19use ordered_float::NotNan;
20use snafu::{OptionExt, ResultExt};
21use vrl::prelude::Bytes;
22use vrl::value::{KeyString, Value as VrlValue};
23
24use crate::error::{
25    FloatIsNanSnafu, Result, ValueDefaultValueUnsupportedSnafu, ValueInvalidResolutionSnafu,
26    ValueParseBooleanSnafu, ValueParseFloatSnafu, ValueParseIntSnafu, ValueParseTypeSnafu,
27    ValueUnsupportedYamlTypeSnafu, ValueYamlKeyMustBeStringSnafu,
28};
29
30pub(crate) const NANOSECOND_RESOLUTION: &str = "nanosecond";
31pub(crate) const NANO_RESOLUTION: &str = "nano";
32pub(crate) const NS_RESOLUTION: &str = "ns";
33pub(crate) const MICROSECOND_RESOLUTION: &str = "microsecond";
34pub(crate) const MICRO_RESOLUTION: &str = "micro";
35pub(crate) const US_RESOLUTION: &str = "us";
36pub(crate) const MILLISECOND_RESOLUTION: &str = "millisecond";
37pub(crate) const MILLI_RESOLUTION: &str = "milli";
38pub(crate) const MS_RESOLUTION: &str = "ms";
39pub(crate) const SECOND_RESOLUTION: &str = "second";
40pub(crate) const SEC_RESOLUTION: &str = "sec";
41pub(crate) const S_RESOLUTION: &str = "s";
42
43pub(crate) const VALID_RESOLUTIONS: [&str; 12] = [
44    NANOSECOND_RESOLUTION,
45    NANO_RESOLUTION,
46    NS_RESOLUTION,
47    MICROSECOND_RESOLUTION,
48    MICRO_RESOLUTION,
49    US_RESOLUTION,
50    MILLISECOND_RESOLUTION,
51    MILLI_RESOLUTION,
52    MS_RESOLUTION,
53    SECOND_RESOLUTION,
54    SEC_RESOLUTION,
55    S_RESOLUTION,
56];
57
58pub fn parse_str_type(t: &str) -> Result<ColumnDataType> {
59    let mut parts = t.splitn(2, ',');
60    let head = parts.next().unwrap_or_default();
61    let tail = parts.next().map(|s| s.trim().to_string());
62    match head.to_lowercase().as_str() {
63        "int8" => Ok(ColumnDataType::Int8),
64        "int16" => Ok(ColumnDataType::Int16),
65        "int32" => Ok(ColumnDataType::Int32),
66        "int64" => Ok(ColumnDataType::Int64),
67
68        "uint8" => Ok(ColumnDataType::Uint8),
69        "uint16" => Ok(ColumnDataType::Uint16),
70        "uint32" => Ok(ColumnDataType::Uint32),
71        "uint64" => Ok(ColumnDataType::Uint64),
72
73        "float32" => Ok(ColumnDataType::Float32),
74        "float64" => Ok(ColumnDataType::Float64),
75
76        "boolean" => Ok(ColumnDataType::Boolean),
77        "string" => Ok(ColumnDataType::String),
78
79        "timestamp" | "epoch" | "time" => match tail {
80            Some(resolution) if !resolution.is_empty() => match resolution.as_str() {
81                NANOSECOND_RESOLUTION | NANO_RESOLUTION | NS_RESOLUTION => {
82                    Ok(ColumnDataType::TimestampNanosecond)
83                }
84                MICROSECOND_RESOLUTION | MICRO_RESOLUTION | US_RESOLUTION => {
85                    Ok(ColumnDataType::TimestampMicrosecond)
86                }
87                MILLISECOND_RESOLUTION | MILLI_RESOLUTION | MS_RESOLUTION => {
88                    Ok(ColumnDataType::TimestampMillisecond)
89                }
90                SECOND_RESOLUTION | SEC_RESOLUTION | S_RESOLUTION => {
91                    Ok(ColumnDataType::TimestampSecond)
92                }
93                _ => ValueInvalidResolutionSnafu {
94                    resolution,
95                    valid_resolution: VALID_RESOLUTIONS.join(","),
96                }
97                .fail(),
98            },
99            _ => Ok(ColumnDataType::TimestampNanosecond),
100        },
101
102        // We only consider object and array to be json types. and use Map to represent json
103        // TODO(qtang): Needs to be defined with better semantics
104        "json" => Ok(ColumnDataType::Binary),
105        "json2" => Ok(ColumnDataType::Json),
106
107        _ => ValueParseTypeSnafu { t }.fail(),
108    }
109}
110
111pub fn parse_str_value(type_: &ColumnDataType, v: &str) -> Result<ValueData> {
112    match type_ {
113        ColumnDataType::Int8 => v
114            .parse::<i8>()
115            .map(|v| ValueData::I8Value(v as i32))
116            .context(ValueParseIntSnafu { ty: "int8", v }),
117        ColumnDataType::Int16 => v
118            .parse::<i16>()
119            .map(|v| ValueData::I16Value(v as i32))
120            .context(ValueParseIntSnafu { ty: "int16", v }),
121        ColumnDataType::Int32 => v
122            .parse::<i32>()
123            .map(ValueData::I32Value)
124            .context(ValueParseIntSnafu { ty: "int32", v }),
125        ColumnDataType::Int64 => v
126            .parse::<i64>()
127            .map(ValueData::I64Value)
128            .context(ValueParseIntSnafu { ty: "int64", v }),
129
130        ColumnDataType::Uint8 => v
131            .parse::<u8>()
132            .map(|v| ValueData::U8Value(v as u32))
133            .context(ValueParseIntSnafu { ty: "uint8", v }),
134        ColumnDataType::Uint16 => v
135            .parse::<u16>()
136            .map(|v| ValueData::U16Value(v as u32))
137            .context(ValueParseIntSnafu { ty: "uint16", v }),
138        ColumnDataType::Uint32 => v
139            .parse::<u32>()
140            .map(ValueData::U32Value)
141            .context(ValueParseIntSnafu { ty: "uint32", v }),
142        ColumnDataType::Uint64 => v
143            .parse::<u64>()
144            .map(ValueData::U64Value)
145            .context(ValueParseIntSnafu { ty: "uint64", v }),
146
147        ColumnDataType::Float32 => v
148            .parse::<f32>()
149            .map(ValueData::F32Value)
150            .context(ValueParseFloatSnafu { ty: "float32", v }),
151        ColumnDataType::Float64 => v
152            .parse::<f64>()
153            .map(ValueData::F64Value)
154            .context(ValueParseFloatSnafu { ty: "float64", v }),
155
156        ColumnDataType::Boolean => v
157            .parse::<bool>()
158            .map(ValueData::BoolValue)
159            .context(ValueParseBooleanSnafu { ty: "boolean", v }),
160        ColumnDataType::String => Ok(ValueData::StringValue(v.to_string())),
161
162        _ => ValueDefaultValueUnsupportedSnafu {
163            value: format!("{:?}", type_),
164        }
165        .fail(),
166    }
167}
168
169pub fn yaml_to_vrl_value(v: &yaml_rust::Yaml) -> Result<VrlValue> {
170    match v {
171        yaml_rust::Yaml::Null => Ok(VrlValue::Null),
172        yaml_rust::Yaml::Boolean(v) => Ok(VrlValue::Boolean(*v)),
173        yaml_rust::Yaml::Integer(v) => Ok(VrlValue::Integer(*v)),
174        yaml_rust::Yaml::Real(v) => {
175            let f = v
176                .parse::<f64>()
177                .context(ValueParseFloatSnafu { ty: "float64", v })?;
178            NotNan::new(f).map(VrlValue::Float).context(FloatIsNanSnafu)
179        }
180        yaml_rust::Yaml::String(v) => Ok(VrlValue::Bytes(Bytes::from(v.clone()))),
181        yaml_rust::Yaml::Array(arr) => {
182            let mut values = vec![];
183            for v in arr {
184                values.push(yaml_to_vrl_value(v)?);
185            }
186            Ok(VrlValue::Array(values))
187        }
188        yaml_rust::Yaml::Hash(v) => {
189            let mut values = BTreeMap::new();
190            for (k, v) in v {
191                let key = k
192                    .as_str()
193                    .with_context(|| ValueYamlKeyMustBeStringSnafu { value: v.clone() })?;
194                values.insert(KeyString::from(key), yaml_to_vrl_value(v)?);
195            }
196            Ok(VrlValue::Object(values))
197        }
198        _ => ValueUnsupportedYamlTypeSnafu { value: v.clone() }.fail(),
199    }
200}