-
Notifications
You must be signed in to change notification settings - Fork 2.1k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
coprocessor/dag: use expression in dag #2261
Merged
Merged
Changes from all commits
Commits
Show all changes
31 commits
Select commit
Hold shift + click to select a range
aff0c62
expr/builtin_cast: implement eval for expression
AndreMouche 71035a2
Merge branch 'master' into shirly/dag_expr
AndreMouche 4cffa6d
expr/builtin_cast: remove offset check in build
AndreMouche 8a583bc
expr/builtin_cast: merge master && fix bug
AndreMouche 5b4c651
expr/mod.rs: address comments
AndreMouche 2662df1
expr/mod.rs: use new expression in selection
AndreMouche 0c080df
dag/executor: use expression in aggregation
AndreMouche 4f24fd7
expr/mod.rs: use new expression in topn
AndreMouche 7187e19
merge master && fix conflicts
AndreMouche 877aced
merge master
AndreMouche 38ae0bd
dag/executor: refactor
AndreMouche 93d4005
Merge branch 'shirly/dag_expression' of github.com:pingcap/tikv into …
AndreMouche 719a18e
Merge branch 'master' into shirly/dag_expression
AndreMouche b2a7065
Merge branch 'master' into shirly/dag_expression
AndreMouche 64f6b6b
Merge branch 'master' into shirly/dag_expression
AndreMouche 60a0694
Merge branch 'master' into shirly/dag_expression
AndreMouche 9adaaf7
Merge branch 'master' into shirly/dag_expression
AndreMouche 88bd51f
Merge branch 'master' into shirly/dag_expression
AndreMouche 3ffee7e
executor/*: address comments
AndreMouche 134f3bd
Merge branch 'master' into shirly/dag_expression
AndreMouche 95e043b
Merge branch 'master' into shirly/dag_expression
AndreMouche a734f13
dag/expr: remove unuseful code
AndreMouche 3a1d004
executor/aggregation: address comments
AndreMouche 3e9b604
Merge branch 'master' into shirly/dag_expression
AndreMouche 09c0fa2
executor/aggregation: address comments
AndreMouche 2355555
dag/*: address comments
AndreMouche 465ef5b
Merge branch 'master' into shirly/dag_expression
AndreMouche cc2bb05
Merge branch 'master' into shirly/dag_expression
AndreMouche e68aaa9
Merge branch 'master' into shirly/dag_expression
AndreMouche 366c9c7
Merge branch 'master' into shirly/dag_expression
AndreMouche 444a357
Merge branch 'master' into shirly/dag_expression
AndreMouche File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -15,22 +15,66 @@ use std::rc::Rc; | |
|
||
use tipb::schema::ColumnInfo; | ||
use tipb::executor::Aggregation; | ||
use tipb::expression::Expr; | ||
use tipb::expression::{Expr, ExprType}; | ||
use util::collections::{HashMap, HashMapEntry as Entry}; | ||
|
||
use coprocessor::codec::table::RowColsDict; | ||
use coprocessor::codec::datum::{self, approximate_size, Datum, DatumEncoder}; | ||
use coprocessor::endpoint::SINGLE_GROUP; | ||
use coprocessor::select::aggregate::{self, AggrFunc}; | ||
use coprocessor::select::xeval::{EvalContext, Evaluator}; | ||
use coprocessor::select::xeval::EvalContext; | ||
use coprocessor::dag::expr::Expression; | ||
use coprocessor::metrics::*; | ||
use coprocessor::Result; | ||
|
||
use super::{inflate_with_col_for_dag, Executor, ExprColumnRefVisitor, Row}; | ||
|
||
struct AggrFuncExpr { | ||
args: Vec<Expression>, | ||
tp: ExprType, | ||
} | ||
|
||
impl AggrFuncExpr { | ||
fn batch_build(ctx: &EvalContext, expr: Vec<Expr>) -> Result<Vec<AggrFuncExpr>> { | ||
let res: Vec<AggrFuncExpr> = try!( | ||
expr.into_iter() | ||
.map(|v| AggrFuncExpr::build(ctx, v)) | ||
.collect() | ||
); | ||
Ok(res) | ||
} | ||
|
||
fn build(ctx: &EvalContext, mut expr: Expr) -> Result<AggrFuncExpr> { | ||
let args = box_try!(Expression::batch_build( | ||
ctx, | ||
expr.take_children().into_vec() | ||
)); | ||
let tp = expr.get_tp(); | ||
Ok(AggrFuncExpr { args: args, tp: tp }) | ||
} | ||
|
||
fn eval_args(&self, ctx: &EvalContext, row: &[Datum]) -> Result<Vec<Datum>> { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. s/eval_args/eval/ |
||
let res: Vec<Datum> = box_try!(self.args.iter().map(|v| v.eval(ctx, row)).collect()); | ||
Ok(res) | ||
} | ||
} | ||
|
||
impl AggrFunc { | ||
fn update_with_expr( | ||
&mut self, | ||
ctx: &EvalContext, | ||
expr: &AggrFuncExpr, | ||
row: &[Datum], | ||
) -> Result<()> { | ||
let vals = try!(expr.eval_args(ctx, row)); | ||
try!(self.update(ctx, vals)); | ||
Ok(()) | ||
} | ||
} | ||
|
||
pub struct AggregationExecutor<'a> { | ||
group_by: Vec<Expr>, | ||
aggr_func: Vec<Expr>, | ||
group_by: Vec<Expression>, | ||
aggr_func: Vec<AggrFuncExpr>, | ||
group_keys: Vec<Rc<Vec<u8>>>, | ||
group_key_aggrs: HashMap<Rc<Vec<u8>>, Vec<Box<AggrFunc>>>, | ||
cursor: usize, | ||
|
@@ -58,8 +102,8 @@ impl<'a> AggregationExecutor<'a> { | |
.with_label_values(&["aggregation"]) | ||
.inc(); | ||
Ok(AggregationExecutor { | ||
group_by: group_by, | ||
aggr_func: aggr_func, | ||
group_by: box_try!(Expression::batch_build(ctx.as_ref(), group_by)), | ||
aggr_func: try!(AggrFuncExpr::batch_build(ctx.as_ref(), aggr_func)), | ||
group_keys: vec![], | ||
group_key_aggrs: map![], | ||
cursor: 0, | ||
|
@@ -71,14 +115,14 @@ impl<'a> AggregationExecutor<'a> { | |
}) | ||
} | ||
|
||
fn get_group_key(&mut self, eval: &mut Evaluator) -> Result<Vec<u8>> { | ||
fn get_group_key(&self, row: &[Datum]) -> Result<Vec<u8>> { | ||
if self.group_by.is_empty() { | ||
let single_group = Datum::Bytes(SINGLE_GROUP.to_vec()); | ||
return Ok(box_try!(datum::encode_value(&[single_group]))); | ||
} | ||
let mut vals = Vec::with_capacity(self.group_by.len()); | ||
for expr in &self.group_by { | ||
let v = box_try!(eval.eval(&self.ctx, expr)); | ||
let v = box_try!(expr.eval(&self.ctx, row)); | ||
vals.push(v); | ||
} | ||
let res = box_try!(datum::encode_value(&vals)); | ||
|
@@ -87,23 +131,20 @@ impl<'a> AggregationExecutor<'a> { | |
|
||
fn aggregate(&mut self) -> Result<()> { | ||
while let Some(row) = try!(self.src.next()) { | ||
let mut eval = Evaluator::default(); | ||
try!(inflate_with_col_for_dag( | ||
&mut eval, | ||
let cols = try!(inflate_with_col_for_dag( | ||
&self.ctx, | ||
&row.data, | ||
self.cols.clone(), | ||
&self.related_cols_offset, | ||
row.handle | ||
)); | ||
let group_key = Rc::new(try!(self.get_group_key(&mut eval))); | ||
let group_key = Rc::new(try!(self.get_group_key(&cols))); | ||
match self.group_key_aggrs.entry(group_key.clone()) { | ||
Entry::Vacant(e) => { | ||
let mut aggrs = Vec::with_capacity(self.aggr_func.len()); | ||
for expr in &self.aggr_func { | ||
let mut aggr = try!(aggregate::build_aggr_func(expr)); | ||
let vals = box_try!(eval.batch_eval(&self.ctx, expr.get_children())); | ||
try!(aggr.update(&self.ctx, vals)); | ||
let mut aggr = try!(aggregate::build_aggr_func(expr.tp)); | ||
try!(aggr.update_with_expr(&self.ctx, expr, &cols)); | ||
aggrs.push(aggr); | ||
} | ||
self.group_keys.push(group_key); | ||
|
@@ -112,8 +153,7 @@ impl<'a> AggregationExecutor<'a> { | |
Entry::Occupied(e) => { | ||
let aggrs = e.into_mut(); | ||
for (expr, aggr) in self.aggr_func.iter().zip(aggrs) { | ||
let vals = box_try!(eval.batch_eval(&self.ctx, expr.get_children())); | ||
box_try!(aggr.update(&self.ctx, vals)); | ||
try!(aggr.update_with_expr(&self.ctx, expr, &cols)); | ||
} | ||
} | ||
} | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Can this be merged with
AggrFunc
?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
address comments