forked from apache/horaedb
-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathread.rs
220 lines (196 loc) · 7.47 KB
/
read.rs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
// Copyright 2023 CeresDB Project Authors. Licensed under Apache-2.0.
//! Contains common methods used by the read process.
use std::time::Instant;
use ceresdbproto::storage::{
storage_service_client::StorageServiceClient, RequestContext, SqlQueryRequest, SqlQueryResponse,
};
use common_types::request_id::RequestId;
use common_util::{error::BoxError, time::InstantExt};
use futures::FutureExt;
use http::StatusCode;
use interpreters::interpreter::Output;
use log::{error, info, warn};
use query_engine::executor::Executor as QueryExecutor;
use query_frontend::{
frontend,
frontend::{Context as SqlContext, Frontend},
provider::CatalogMetaProvider,
};
use router::endpoint::Endpoint;
use snafu::{ensure, ResultExt};
use tonic::{transport::Channel, IntoRequest};
use crate::{
error::{ErrNoCause, ErrWithCause, Error, Internal, Result},
forward::{ForwardRequest, ForwardResult},
Context, Proxy,
};
pub enum SqlResponse {
Forwarded(SqlQueryResponse),
Local(Output),
}
impl<Q: QueryExecutor + 'static> Proxy<Q> {
pub(crate) async fn handle_sql(
&self,
ctx: Context,
schema: &str,
sql: &str,
) -> Result<SqlResponse> {
if let Some(resp) = self
.maybe_forward_sql_query(ctx.clone(), schema, sql)
.await?
{
match resp {
ForwardResult::Forwarded(resp) => return Ok(SqlResponse::Forwarded(resp?)),
ForwardResult::Local => (),
}
};
Ok(SqlResponse::Local(
self.fetch_sql_query_output(ctx, schema, sql).await?,
))
}
pub(crate) async fn fetch_sql_query_output(
&self,
ctx: Context,
schema: &str,
sql: &str,
) -> Result<Output> {
let request_id = RequestId::next_id();
let begin_instant = Instant::now();
let deadline = ctx.timeout.map(|t| begin_instant + t);
let catalog = self.instance.catalog_manager.default_catalog_name();
info!("Handle sql query, request_id:{request_id}, schema:{schema}, sql:{sql}");
let instance = &self.instance;
// TODO(yingwen): Privilege check, cannot access data of other tenant
// TODO(yingwen): Maybe move MetaProvider to instance
let provider = CatalogMetaProvider {
manager: instance.catalog_manager.clone(),
default_catalog: catalog,
default_schema: schema,
function_registry: &*instance.function_registry,
};
let frontend = Frontend::new(provider);
let mut sql_ctx = SqlContext::new(request_id, deadline);
// Parse sql, frontend error of invalid sql already contains sql
// TODO(yingwen): Maybe move sql from frontend error to outer error
let mut stmts = frontend
.parse_sql(&mut sql_ctx, sql)
.box_err()
.context(ErrWithCause {
code: StatusCode::BAD_REQUEST,
msg: "Failed to parse sql",
})?;
ensure!(
!stmts.is_empty(),
ErrNoCause {
code: StatusCode::BAD_REQUEST,
msg: format!("No valid query statement provided, sql:{sql}",),
}
);
// TODO(yingwen): For simplicity, we only support executing one statement now
// TODO(yingwen): INSERT/UPDATE/DELETE can be batched
ensure!(
stmts.len() == 1,
ErrNoCause {
code: StatusCode::BAD_REQUEST,
msg: format!(
"Only support execute one statement now, current num:{}, sql:{}",
stmts.len(),
sql
),
}
);
// Open partition table if needed.
let table_name = frontend::parse_table_name(&stmts);
if let Some(table_name) = &table_name {
self.maybe_open_partition_table_if_not_exist(catalog, schema, table_name)
.await?;
}
// Create logical plan
// Note: Remember to store sql in error when creating logical plan
let plan = frontend
// TODO(yingwen): Check error, some error may indicate that the sql is invalid. Now we
// return internal server error in those cases
.statement_to_plan(&mut sql_ctx, stmts.remove(0))
.box_err()
.with_context(|| ErrWithCause {
code: StatusCode::INTERNAL_SERVER_ERROR,
msg: format!("Failed to create plan, query:{sql}"),
})?;
if let Some(table_name) = &table_name {
if !self.valid_ttl_range(&plan, catalog, schema, table_name) {
return Err(Error::SqlQueryOverTTL {
code: StatusCode::OK,
msg: format!("The query {sql}'s timestamp is already out of date"),
});
}
}
let output = if ctx.enable_partition_table_access {
self.execute_plan_involving_partition_table(request_id, catalog, schema, plan, deadline)
.await
} else {
self.execute_plan(request_id, catalog, schema, plan, deadline)
.await
};
let output = output.box_err().with_context(|| ErrWithCause {
code: StatusCode::INTERNAL_SERVER_ERROR,
msg: format!("Failed to execute plan, sql:{sql}"),
})?;
let cost = begin_instant.saturating_elapsed();
info!("Handle sql query success, catalog:{catalog}, schema:{schema}, request_id:{request_id}, cost:{cost:?}, sql:{sql:?}");
Ok(output)
}
async fn maybe_forward_sql_query(
&self,
ctx: Context,
schema: &str,
sql: &str,
) -> Result<Option<ForwardResult<SqlQueryResponse, Error>>> {
let table_name = frontend::parse_table_name_with_sql(sql)
.box_err()
.with_context(|| Internal {
msg: format!("Failed to parse table name with sql, sql:{sql}"),
})?;
if table_name.is_none() {
warn!("Unable to forward sql query without table name, sql:{sql}",);
return Ok(None);
}
let sql_request = SqlQueryRequest {
context: Some(RequestContext {
database: schema.to_string(),
}),
tables: vec![],
sql: sql.to_string(),
};
let forward_req = ForwardRequest {
schema: schema.to_string(),
table: table_name.unwrap(),
req: sql_request.into_request(),
forwarded_from: ctx.forwarded_from,
};
let do_query = |mut client: StorageServiceClient<Channel>,
request: tonic::Request<SqlQueryRequest>,
_: &Endpoint| {
let query = async move {
client
.sql_query(request)
.await
.map(|resp| resp.into_inner())
.box_err()
.context(ErrWithCause {
code: StatusCode::INTERNAL_SERVER_ERROR,
msg: "Forwarded sql query failed",
})
}
.boxed();
Box::new(query) as _
};
let forward_result = self.forwarder.forward(forward_req, do_query).await;
Ok(match forward_result {
Ok(forward_res) => Some(forward_res),
Err(e) => {
error!("Failed to forward sql req but the error is ignored, err:{e}");
None
}
})
}
}