Skip to main content

frontend/instance/
jaeger.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, HashSet};
16use std::sync::Arc;
17
18use async_trait::async_trait;
19use auth::{JAEGER_QUERY, PermissionReq, PermissionTableTarget, PermissionTableTargets};
20use catalog::CatalogManagerRef;
21use common_catalog::consts::{
22    TRACE_TABLE_NAME, trace_operations_table_name, trace_services_table_name,
23};
24use common_function::function::FunctionRef;
25use common_function::scalars::json::json_get::{
26    JsonGetBool, JsonGetFloat, JsonGetInt, JsonGetString, JsonGetWithType,
27};
28use common_function::scalars::udf::create_udf;
29use common_query::{Output, OutputData};
30use common_recordbatch::adapter::RecordBatchStreamAdapter;
31use common_recordbatch::util;
32use common_telemetry::warn;
33use datafusion::common::ScalarValue;
34use datafusion::dataframe::DataFrame;
35use datafusion::execution::SessionStateBuilder;
36use datafusion::execution::context::SessionContext;
37use datafusion::functions::core::expr_fn::coalesce;
38use datafusion::functions_window::expr_fn::row_number;
39use datafusion_expr::select_expr::SelectExpr;
40use datafusion_expr::{Expr, ExprFunctionExt, SortExpr, col, lit, lit_timestamp_nano, wildcard};
41use query::QueryEngineRef;
42use serde_json::Value as JsonValue;
43use servers::error::{
44    AuthSnafu, CollectRecordbatchSnafu, DataFusionSnafu, Result as ServerResult, TableNotFoundSnafu,
45};
46use servers::http::jaeger::{JAEGER_QUERY_TABLE_NAME_KEY, QueryTraceParams, TraceUserAgent};
47use servers::otlp::trace::{
48    DURATION_NANO_COLUMN, KEY_OTEL_STATUS_ERROR_KEY, RESOURCE_ATTRIBUTES_COLUMN,
49    SERVICE_NAME_COLUMN, SPAN_ATTRIBUTES_COLUMN, SPAN_KIND_COLUMN, SPAN_KIND_PREFIX,
50    SPAN_NAME_COLUMN, SPAN_STATUS_CODE, SPAN_STATUS_ERROR, TIMESTAMP_COLUMN, TRACE_ID_COLUMN,
51};
52use servers::query_handler::JaegerQueryHandler;
53use session::context::QueryContextRef;
54use snafu::{OptionExt, ResultExt};
55use table::TableRef;
56use table::table::adapter::DfTableProviderAdapter;
57
58use crate::instance::Instance;
59
60const DEFAULT_LIMIT: usize = 2000;
61const KEY_RN: &str = "greptime_rn";
62
63impl Instance {
64    async fn check_jaeger_query_permission(&self, ctx: &QueryContextRef) -> ServerResult<()> {
65        let table = ctx
66            .extension(JAEGER_QUERY_TABLE_NAME_KEY)
67            .unwrap_or(TRACE_TABLE_NAME);
68        let targets = PermissionTableTargets::resolved(vec![PermissionTableTarget::new(
69            ctx.current_catalog(),
70            ctx.current_schema(),
71            table,
72        )]);
73        let targets = self.resolve_query_permission_targets(targets, ctx).await?;
74        self.check_table_permission(ctx, PermissionReq::Action(JAEGER_QUERY), targets)
75            .context(AuthSnafu)?;
76        Ok(())
77    }
78}
79
80#[async_trait]
81impl JaegerQueryHandler for Instance {
82    async fn get_services(&self, ctx: QueryContextRef) -> ServerResult<Output> {
83        self.check_jaeger_query_permission(&ctx).await?;
84
85        // It's equivalent to `SELECT DISTINCT(service_name) FROM {db}.{trace_table}`.
86        Ok(query_trace_table(
87            ctx,
88            self,
89            vec![SelectExpr::from(col(SERVICE_NAME_COLUMN))],
90            vec![],
91            vec![],
92            None,
93            None,
94            vec![col(SERVICE_NAME_COLUMN)],
95        )
96        .await?)
97    }
98
99    async fn get_operations(
100        &self,
101        ctx: QueryContextRef,
102        service_name: &str,
103        span_kind: Option<&str>,
104    ) -> ServerResult<Output> {
105        self.check_jaeger_query_permission(&ctx).await?;
106
107        let mut filters = vec![col(SERVICE_NAME_COLUMN).eq(lit(service_name))];
108
109        if let Some(span_kind) = span_kind {
110            filters.push(col(SPAN_KIND_COLUMN).eq(lit(format!(
111                "{}{}",
112                SPAN_KIND_PREFIX,
113                span_kind.to_uppercase()
114            ))));
115        }
116
117        // It's equivalent to the following SQL query:
118        //
119        // ```
120        // SELECT DISTINCT span_name, span_kind
121        // FROM
122        //   {db}.{trace_table}
123        // WHERE
124        //   service_name = '{service_name}' AND
125        //   span_kind = '{span_kind}'
126        // ORDER BY
127        //   span_name ASC
128        // ```.
129        Ok(query_trace_table(
130            ctx,
131            self,
132            vec![
133                SelectExpr::from(col(SPAN_NAME_COLUMN)),
134                SelectExpr::from(col(SPAN_KIND_COLUMN)),
135                SelectExpr::from(col(SERVICE_NAME_COLUMN)),
136                SelectExpr::from(col(TIMESTAMP_COLUMN)),
137            ],
138            filters,
139            vec![col(SPAN_NAME_COLUMN).sort(true, false)], // Sort by span_name in ascending order.
140            Some(DEFAULT_LIMIT),
141            None,
142            vec![col(SPAN_NAME_COLUMN), col(SPAN_KIND_COLUMN)],
143        )
144        .await?)
145    }
146
147    async fn get_trace(
148        &self,
149        ctx: QueryContextRef,
150        trace_id: &str,
151        start_time: Option<i64>,
152        end_time: Option<i64>,
153        limit: Option<usize>,
154    ) -> ServerResult<Output> {
155        self.check_jaeger_query_permission(&ctx).await?;
156
157        // It's equivalent to the following SQL query:
158        //
159        // ```
160        // SELECT
161        //   *
162        // FROM
163        //   {db}.{trace_table}
164        // WHERE
165        //   trace_id = '{trace_id}' AND
166        //   timestamp >= {start_time} AND
167        //   timestamp <= {end_time}
168        // ORDER BY
169        //   timestamp DESC
170        // ```.
171        let selects = vec![wildcard()];
172
173        let mut filters = vec![col(TRACE_ID_COLUMN).eq(lit(trace_id))];
174
175        if let Some(start_time) = start_time {
176            filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
177        }
178
179        if let Some(end_time) = end_time {
180            filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
181        }
182
183        Ok(query_trace_table(
184            ctx,
185            self,
186            selects,
187            filters,
188            vec![col(TIMESTAMP_COLUMN).sort(false, false)], // Sort by timestamp in descending order.
189            limit,
190            None,
191            vec![],
192        )
193        .await?)
194    }
195
196    async fn find_traces(
197        &self,
198        ctx: QueryContextRef,
199        query_params: QueryTraceParams,
200    ) -> ServerResult<Output> {
201        self.check_jaeger_query_permission(&ctx).await?;
202
203        let mut filters = vec![];
204
205        // `service_name` is already validated in `from_jaeger_query_params()`, so no additional check needed here.
206        filters.push(col(SERVICE_NAME_COLUMN).eq(lit(query_params.service_name)));
207
208        if let Some(operation_name) = query_params.operation_name {
209            filters.push(col(SPAN_NAME_COLUMN).eq(lit(operation_name)));
210        }
211
212        if let Some(start_time) = query_params.start_time {
213            filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
214        }
215
216        if let Some(end_time) = query_params.end_time {
217            filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
218        }
219
220        if let Some(min_duration) = query_params.min_duration {
221            filters.push(col(DURATION_NANO_COLUMN).gt_eq(lit(min_duration)));
222        }
223
224        if let Some(max_duration) = query_params.max_duration {
225            filters.push(col(DURATION_NANO_COLUMN).lt_eq(lit(max_duration)));
226        }
227
228        // Get all distinct trace ids that match the filters.
229        // It's equivalent to the following SQL query:
230        //
231        // ```
232        // SELECT DISTINCT trace_id
233        // FROM
234        //   {db}.{trace_table}
235        // WHERE
236        //   service_name = '{service_name}' AND
237        //   operation_name = '{operation_name}' AND
238        //   timestamp >= {start_time} AND
239        //   timestamp <= {end_time} AND
240        //   duration >= {min_duration} AND
241        //   duration <= {max_duration}
242        // LIMIT {limit}
243        // ```.
244        let output = query_trace_table(
245            ctx.clone(),
246            self,
247            vec![wildcard()],
248            filters,
249            vec![],
250            Some(query_params.limit.unwrap_or(DEFAULT_LIMIT)),
251            query_params.tags,
252            vec![col(TRACE_ID_COLUMN)],
253        )
254        .await?;
255
256        // Get all traces that match the trace ids from the previous query.
257        // It's equivalent to the following SQL query:
258        //
259        // ```
260        // SELECT *
261        // FROM
262        //   {db}.{trace_table}
263        // WHERE
264        //   trace_id IN ({trace_ids}) AND
265        //   timestamp >= {start_time} AND
266        //   timestamp <= {end_time}
267        // ```
268        let mut filters = vec![
269            col(TRACE_ID_COLUMN).in_list(
270                trace_ids_from_output(output)
271                    .await?
272                    .iter()
273                    .map(lit)
274                    .collect::<Vec<Expr>>(),
275                false,
276            ),
277        ];
278
279        if let Some(start_time) = query_params.start_time {
280            filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
281        }
282
283        if let Some(end_time) = query_params.end_time {
284            filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
285        }
286
287        match query_params.user_agent {
288            TraceUserAgent::Grafana => {
289                // grafana only use trace id and timestamp
290                // clicking the trace id will invoke the query trace api
291                // so we only need to return 1 span for each trace
292                let table_name = ctx
293                    .extension(JAEGER_QUERY_TABLE_NAME_KEY)
294                    .unwrap_or(TRACE_TABLE_NAME);
295
296                let table = get_table(ctx.clone(), self.catalog_manager(), table_name).await?;
297
298                Ok(find_traces_rank_3(
299                    table,
300                    self.query_engine(),
301                    filters,
302                    vec![col(TIMESTAMP_COLUMN).sort(false, false)], // Sort by timestamp in descending order.
303                )
304                .await?)
305            }
306            _ => {
307                // query all spans
308                Ok(query_trace_table(
309                    ctx,
310                    self,
311                    vec![wildcard()],
312                    filters,
313                    vec![col(TIMESTAMP_COLUMN).sort(false, false)], // Sort by timestamp in descending order.
314                    None,
315                    None,
316                    vec![],
317                )
318                .await?)
319            }
320        }
321    }
322}
323
324#[allow(clippy::too_many_arguments)]
325async fn query_trace_table(
326    ctx: QueryContextRef,
327    instance: &Instance,
328    selects: Vec<SelectExpr>,
329    filters: Vec<Expr>,
330    sorts: Vec<SortExpr>,
331    limit: Option<usize>,
332    tags: Option<HashMap<String, JsonValue>>,
333    distincts: Vec<Expr>,
334) -> ServerResult<Output> {
335    let trace_table_name = ctx
336        .extension(JAEGER_QUERY_TABLE_NAME_KEY)
337        .unwrap_or(TRACE_TABLE_NAME);
338
339    // If only select services, use the trace services table.
340    // If querying operations (distinct by span_name and span_kind), use the trace operations table.
341    let table_name = {
342        if match selects.as_slice() {
343            [SelectExpr::Expression(x)] => x == &col(SERVICE_NAME_COLUMN),
344            _ => false,
345        } {
346            &trace_services_table_name(trace_table_name)
347        } else if !distincts.is_empty()
348            && distincts.contains(&col(SPAN_NAME_COLUMN))
349            && distincts.contains(&col(SPAN_KIND_COLUMN))
350        {
351            &trace_operations_table_name(trace_table_name)
352        } else {
353            trace_table_name
354        }
355    };
356
357    let table = instance
358        .catalog_manager()
359        .table(
360            ctx.current_catalog(),
361            &ctx.current_schema(),
362            table_name,
363            Some(&ctx),
364        )
365        .await?
366        .with_context(|| TableNotFoundSnafu {
367            table: table_name,
368            catalog: ctx.current_catalog(),
369            schema: ctx.current_schema(),
370        })?;
371
372    let table_info = table.table_info();
373    let data_model = table_info.meta.options.data_model();
374
375    // collect to set
376    let col_names = table_info
377        .meta
378        .field_column_names()
379        .map(|s| format!("\"{}\"", s))
380        .collect::<HashSet<String>>();
381
382    let df_context = create_df_context(instance.query_engine())?;
383
384    let dataframe = df_context
385        .read_table(Arc::new(DfTableProviderAdapter::new(table)))
386        .context(DataFusionSnafu)?;
387
388    let dataframe = dataframe.select(selects).context(DataFusionSnafu)?;
389
390    // Apply all filters.
391    let dataframe = filters
392        .into_iter()
393        .chain(tags.map_or(Ok(vec![]), |t| {
394            tags_filters(&dataframe, t, data_model, &col_names)
395        })?)
396        .try_fold(dataframe, |df, expr| {
397            df.filter(expr).context(DataFusionSnafu)
398        })?;
399
400    // Apply the distinct if needed.
401    let dataframe = if !distincts.is_empty() {
402        dataframe
403            .distinct_on(distincts.clone(), distincts, None)
404            .context(DataFusionSnafu)?
405    } else {
406        dataframe
407    };
408
409    // Apply the sorts if needed.
410    let dataframe = if !sorts.is_empty() {
411        dataframe.sort(sorts).context(DataFusionSnafu)?
412    } else {
413        dataframe
414    };
415
416    // Apply the limit if needed.
417    let dataframe = if let Some(limit) = limit {
418        dataframe.limit(0, Some(limit)).context(DataFusionSnafu)?
419    } else {
420        dataframe
421    };
422
423    // Execute the query and collect the result.
424    let stream = dataframe.execute_stream().await.context(DataFusionSnafu)?;
425
426    let output = Output::new_with_stream(Box::pin(
427        RecordBatchStreamAdapter::try_new(stream).context(CollectRecordbatchSnafu)?,
428    ));
429
430    output
431        .map_dictionary_to_values()
432        .context(CollectRecordbatchSnafu)
433}
434
435async fn get_table(
436    ctx: QueryContextRef,
437    catalog_manager: &CatalogManagerRef,
438    table_name: &str,
439) -> ServerResult<TableRef> {
440    catalog_manager
441        .table(
442            ctx.current_catalog(),
443            &ctx.current_schema(),
444            table_name,
445            Some(&ctx),
446        )
447        .await?
448        .with_context(|| TableNotFoundSnafu {
449            table: table_name,
450            catalog: ctx.current_catalog(),
451            schema: ctx.current_schema(),
452        })
453}
454
455async fn find_traces_rank_3(
456    table: TableRef,
457    query_engine: &QueryEngineRef,
458    filters: Vec<Expr>,
459    sorts: Vec<SortExpr>,
460) -> ServerResult<Output> {
461    let df_context = create_df_context(query_engine)?;
462
463    let dataframe = df_context
464        .read_table(Arc::new(DfTableProviderAdapter::new(table)))
465        .context(DataFusionSnafu)?;
466
467    let dataframe = dataframe
468        .select(vec![wildcard()])
469        .context(DataFusionSnafu)?;
470
471    // Apply all filters.
472    let dataframe = filters.into_iter().try_fold(dataframe, |df, expr| {
473        df.filter(expr).context(DataFusionSnafu)
474    })?;
475
476    // Apply the sorts if needed.
477    let dataframe = if !sorts.is_empty() {
478        dataframe.sort(sorts).context(DataFusionSnafu)?
479    } else {
480        dataframe
481    };
482
483    // create rank column, for each trace, get the earliest 3 spans
484    let trace_id_col = vec![col(TRACE_ID_COLUMN)];
485    let timestamp_asc = vec![col(TIMESTAMP_COLUMN).sort(true, false)];
486
487    let dataframe = dataframe
488        .with_column(
489            KEY_RN,
490            row_number()
491                .partition_by(trace_id_col)
492                .order_by(timestamp_asc)
493                .build()
494                .context(DataFusionSnafu)?,
495        )
496        .context(DataFusionSnafu)?;
497
498    let dataframe = dataframe
499        .filter(col(KEY_RN).lt_eq(lit(3)))
500        .context(DataFusionSnafu)?;
501
502    // Execute the query and collect the result.
503    let stream = dataframe.execute_stream().await.context(DataFusionSnafu)?;
504
505    let output = Output::new_with_stream(Box::pin(
506        RecordBatchStreamAdapter::try_new(stream).context(CollectRecordbatchSnafu)?,
507    ));
508
509    output
510        .map_dictionary_to_values()
511        .context(CollectRecordbatchSnafu)
512}
513
514// The current implementation registers UDFs during the planning stage, which makes it difficult
515// to utilize them through DataFrame APIs. To address this limitation, we create a new session
516// context and register the required UDFs, allowing them to be decoupled from the global context.
517// TODO(zyy17): Is it possible or necessary to reuse the existing session context?
518fn create_df_context(query_engine: &QueryEngineRef) -> ServerResult<SessionContext> {
519    let df_context = SessionContext::new_with_state(
520        SessionStateBuilder::new_from_existing(query_engine.engine_state().session_state()).build(),
521    );
522
523    // JSON UDFs used by the v0 and v2 tag filters.
524    let udfs: Vec<FunctionRef> = vec![
525        Arc::new(JsonGetWithType::default()),
526        Arc::new(JsonGetInt::default()),
527        Arc::new(JsonGetFloat::default()),
528        Arc::new(JsonGetBool::default()),
529        Arc::new(JsonGetString::default()),
530    ];
531
532    for udf in udfs {
533        df_context.register_udf(create_udf(udf));
534    }
535
536    Ok(df_context)
537}
538
539fn json_tag_filters(
540    dataframe: &DataFrame,
541    tags: HashMap<String, JsonValue>,
542) -> ServerResult<Vec<Expr>> {
543    let mut filters = vec![];
544
545    // NOTE: The key of the tags may contain `.`, for example: `http.status_code`, so we need to use `["http.status_code"]` in json path to access the value.
546    for (key, value) in tags.iter() {
547        if let JsonValue::String(value) = value {
548            filters.push(
549                dataframe
550                    .registry()
551                    .udf(JsonGetString::NAME)
552                    .context(DataFusionSnafu)?
553                    .call(vec![
554                        col(SPAN_ATTRIBUTES_COLUMN),
555                        lit(format!("[\"{}\"]", key)),
556                    ])
557                    .eq(lit(value)),
558            );
559        }
560        if let JsonValue::Number(value) = value {
561            if value.is_i64() {
562                filters.push(
563                    dataframe
564                        .registry()
565                        .udf(JsonGetInt::NAME)
566                        .context(DataFusionSnafu)?
567                        .call(vec![
568                            col(SPAN_ATTRIBUTES_COLUMN),
569                            lit(format!("[\"{}\"]", key)),
570                        ])
571                        .eq(lit(value.as_i64().unwrap())),
572                );
573            }
574            if value.is_f64() {
575                filters.push(
576                    dataframe
577                        .registry()
578                        .udf(JsonGetFloat::NAME)
579                        .context(DataFusionSnafu)?
580                        .call(vec![
581                            col(SPAN_ATTRIBUTES_COLUMN),
582                            lit(format!("[\"{}\"]", key)),
583                        ])
584                        .eq(lit(value.as_f64().unwrap())),
585                );
586            }
587        }
588        if let JsonValue::Bool(value) = value {
589            filters.push(
590                dataframe
591                    .registry()
592                    .udf(JsonGetBool::NAME)
593                    .context(DataFusionSnafu)?
594                    .call(vec![
595                        col(SPAN_ATTRIBUTES_COLUMN),
596                        lit(format!("[\"{}\"]", key)),
597                    ])
598                    .eq(lit(*value)),
599            );
600        }
601    }
602
603    Ok(filters)
604}
605
606/// Resolves each tag per row, falling back to Resource when the typed Span value is null.
607fn json2_tag_filters(
608    dataframe: &DataFrame,
609    tags: HashMap<String, JsonValue>,
610) -> ServerResult<Vec<Expr>> {
611    let get = dataframe
612        .registry()
613        .udf(JsonGetWithType::NAME)
614        .context(DataFusionSnafu)?;
615    tags.into_iter()
616        .map(|(key, value)| {
617            if key == KEY_OTEL_STATUS_ERROR_KEY && value == JsonValue::Bool(true) {
618                return Ok(col(SPAN_STATUS_CODE).eq(lit(SPAN_STATUS_ERROR)));
619            }
620            // JSON quoting preserves literal dots, quotes, and backslashes in attribute keys.
621            let path = lit(format!("$[{}]", JsonValue::String(key)));
622            let (value, value_type) = match value {
623                JsonValue::String(value) => (
624                    ScalarValue::Utf8View(Some(value)),
625                    ScalarValue::Utf8View(None),
626                ),
627                JsonValue::Bool(value) => (
628                    ScalarValue::Boolean(Some(value)),
629                    ScalarValue::Boolean(None),
630                ),
631                JsonValue::Number(value) => {
632                    if let Some(value) = value.as_i64() {
633                        (ScalarValue::Int64(Some(value)), ScalarValue::Int64(None))
634                    } else if let Some(value) = value.as_u64() {
635                        (ScalarValue::UInt64(Some(value)), ScalarValue::UInt64(None))
636                    } else if let Some(value) = value.as_f64() {
637                        (
638                            ScalarValue::Float64(Some(value)),
639                            ScalarValue::Float64(None),
640                        )
641                    } else {
642                        return Ok(lit(false));
643                    }
644                }
645                JsonValue::Null => (ScalarValue::Utf8View(None), ScalarValue::Utf8View(None)),
646                JsonValue::Array(_) | JsonValue::Object(_) => return Ok(lit(false)),
647            };
648            let attribute = coalesce(vec![
649                get.call(vec![
650                    col(SPAN_ATTRIBUTES_COLUMN),
651                    path.clone(),
652                    lit(value_type.clone()),
653                ]),
654                get.call(vec![col(RESOURCE_ATTRIBUTES_COLUMN), path, lit(value_type)]),
655            ]);
656            Ok(if value.is_null() {
657                attribute.is_null()
658            } else {
659                attribute.eq(lit(value))
660            })
661        })
662        .collect()
663}
664
665/// Helper function to check if span_key or resource_key exists in col_names and create an expression.
666/// If neither exists, logs a warning and returns None.
667#[inline]
668fn check_col_and_build_expr<F>(
669    span_key: String,
670    resource_key: String,
671    key: &str,
672    col_names: &HashSet<String>,
673    expr_builder: F,
674) -> Option<Expr>
675where
676    F: FnOnce(String) -> Expr,
677{
678    if col_names.contains(&span_key) {
679        return Some(expr_builder(span_key));
680    }
681    if col_names.contains(&resource_key) {
682        return Some(expr_builder(resource_key));
683    }
684    warn!("tag key {} not found in table columns", key);
685    None
686}
687
688fn flatten_tag_filters(
689    tags: HashMap<String, JsonValue>,
690    col_names: &HashSet<String>,
691) -> ServerResult<Vec<Expr>> {
692    let filters = tags
693        .into_iter()
694        .filter_map(|(key, value)| {
695            if key == KEY_OTEL_STATUS_ERROR_KEY && value == JsonValue::Bool(true) {
696                return Some(col(SPAN_STATUS_CODE).eq(lit(SPAN_STATUS_ERROR)));
697            }
698
699            // TODO(shuiyisong): add more precise mapping from key to col name
700            let span_key = format!("\"span_attributes.{}\"", key);
701            let resource_key = format!("\"resource_attributes.{}\"", key);
702            match value {
703                JsonValue::String(value) => {
704                    check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
705                        col(k).eq(lit(value))
706                    })
707                }
708                JsonValue::Number(value) => {
709                    if value.is_f64() {
710                        // safe to unwrap as checked previously
711                        let value = value.as_f64().unwrap();
712                        check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
713                            col(k).eq(lit(value))
714                        })
715                    } else {
716                        let value = value.as_i64().unwrap();
717                        check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
718                            col(k).eq(lit(value))
719                        })
720                    }
721                }
722                JsonValue::Bool(value) => {
723                    check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
724                        col(k).eq(lit(value))
725                    })
726                }
727                JsonValue::Null => {
728                    check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
729                        col(k).is_null()
730                    })
731                }
732                // not supported at the moment
733                JsonValue::Array(_value) => None,
734                JsonValue::Object(_value) => None,
735            }
736        })
737        .collect();
738    Ok(filters)
739}
740
741fn tags_filters(
742    dataframe: &DataFrame,
743    tags: HashMap<String, JsonValue>,
744    data_model: Option<&str>,
745    col_names: &HashSet<String>,
746) -> ServerResult<Vec<Expr>> {
747    match data_model {
748        Some(table::requests::TABLE_DATA_MODEL_TRACE_V1) => flatten_tag_filters(tags, col_names),
749        Some(table::requests::TABLE_DATA_MODEL_TRACE_V2) => json2_tag_filters(dataframe, tags),
750        _ => json_tag_filters(dataframe, tags),
751    }
752}
753
754// Get trace ids from the output in recordbatches.
755async fn trace_ids_from_output(output: Output) -> ServerResult<Vec<String>> {
756    if let OutputData::Stream(stream) = output.data {
757        let schema = stream.schema().clone();
758        let recordbatches = util::collect(stream)
759            .await
760            .context(CollectRecordbatchSnafu)?;
761
762        // Only contains `trace_id` column in string type.
763        if !recordbatches.is_empty()
764            && schema.num_columns() == 1
765            && schema.contains_column(TRACE_ID_COLUMN)
766        {
767            let mut trace_ids = vec![];
768            for recordbatch in recordbatches {
769                recordbatch
770                    .iter_column_as_string(0)
771                    .flatten()
772                    .for_each(|x| trace_ids.push(x));
773            }
774
775            return Ok(trace_ids);
776        }
777    }
778
779    Ok(vec![])
780}