1use 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
45pub 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 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 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 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
243pub 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
256pub 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
265pub 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
284pub 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
352pub(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
436pub 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
454pub 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}