From f84335e9864a3e461bd0d383ce8ae63ffe28d223 Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Wed, 20 May 2026 10:12:20 +0800 Subject: [PATCH] fix(elasticsearch): parse aggregation query results into table format (#344) Previously, aggregation queries (date_histogram, terms, etc.) with `size: 0` returned empty results because the empty hits array was matched but contained no documents, ignoring the aggregations data. --- .../dbx-core/src/db/elasticsearch_driver.rs | 100 +++++++++++++++++- 1 file changed, 99 insertions(+), 1 deletion(-) diff --git a/crates/dbx-core/src/db/elasticsearch_driver.rs b/crates/dbx-core/src/db/elasticsearch_driver.rs index 619309d2b..4406e1f73 100644 --- a/crates/dbx-core/src/db/elasticsearch_driver.rs +++ b/crates/dbx-core/src/db/elasticsearch_driver.rs @@ -249,7 +249,7 @@ pub async fn execute_rest_query(client: &EsClient, input: &str) -> Result::new(); let docs: Vec> = hits .iter() @@ -298,6 +298,31 @@ pub async fn execute_rest_query(client: &EsClient, input: &str) -> Result Result) -> (Vec, Vec>) { + for (_name, agg_value) in aggs { + if let Some(buckets) = agg_value.get("buckets").and_then(|b| b.as_array()) { + if buckets.is_empty() { + continue; + } + let mut all_keys = Vec::::new(); + let mut bucket_rows = Vec::new(); + + for bucket in buckets { + if let Some(obj) = bucket.as_object() { + let mut row = serde_json::Map::new(); + for (k, v) in obj { + if let Some(sub) = v.as_object() { + if let Some(val) = sub.get("value") { + row.insert(k.clone(), val.clone()); + } else { + row.insert(k.clone(), serde_json::Value::String(v.to_string())); + } + } else { + row.insert(k.clone(), v.clone()); + } + } + for key in row.keys() { + if !all_keys.contains(key) { + all_keys.push(key.clone()); + } + } + bucket_rows.push(row); + } + } + + let rows = bucket_rows + .iter() + .map(|br| { + all_keys + .iter() + .map(|k| { + br.get(k) + .map(|v| match v { + serde_json::Value::String(s) => serde_json::Value::String(s.clone()), + other => serde_json::Value::String(other.to_string()), + }) + .unwrap_or(serde_json::Value::Null) + }) + .collect() + }) + .collect(); + + return (all_keys, rows); + } + } + + let mut columns = Vec::new(); + let mut values = Vec::new(); + for (name, agg_value) in aggs { + if let Some(obj) = agg_value.as_object() { + if let Some(val) = obj.get("value") { + columns.push(name.clone()); + values.push(match val { + serde_json::Value::String(s) => serde_json::Value::String(s.clone()), + other => serde_json::Value::String(other.to_string()), + }); + } + } + } + if !columns.is_empty() { + return (columns, vec![values]); + } + + (Vec::new(), Vec::new()) +}