(
&self, mut stream: Stream<Record>, plan: &[pb::PhysicalOpr],
)
| 195 | } |
| 196 | |
| 197 | fn install( |
| 198 | &self, mut stream: Stream<Record>, plan: &[pb::PhysicalOpr], |
| 199 | ) -> Result<Stream<Record>, BuildJobError> { |
| 200 | let mut prev_op_kind = pb::physical_opr::operator::OpKind::Root(pb::Root {}); |
| 201 | for op in &plan[..] { |
| 202 | let op_kind = to_op_kind(op)?; |
| 203 | match op_kind { |
| 204 | OpKind::Repartition(repartition) => { |
| 205 | let repartition_strategy = repartition.strategy.as_ref().ok_or_else(|| { |
| 206 | FnGenError::from(ParsePbError::EmptyFieldError( |
| 207 | "Empty repartition strategy".to_string(), |
| 208 | )) |
| 209 | })?; |
| 210 | match repartition_strategy { |
| 211 | pb::repartition::Strategy::ToAnother(shuffle) => { |
| 212 | let router = self.udf_gen.gen_shuffle(shuffle)?; |
| 213 | stream = stream.repartition(move |t| router.route(t)); |
| 214 | } |
| 215 | pb::repartition::Strategy::ToOthers(_) => stream = stream.broadcast(), |
| 216 | } |
| 217 | } |
| 218 | OpKind::Project(project) => { |
| 219 | let func = self.udf_gen.gen_project(project)?; |
| 220 | stream = stream.filter_map_with_name("Project", move |input| func.exec(input))?; |
| 221 | } |
| 222 | OpKind::Select(select) => { |
| 223 | let func = self.udf_gen.gen_filter(select)?; |
| 224 | stream = stream.filter(move |input| func.test(input))?; |
| 225 | } |
| 226 | OpKind::Unfold(unfold) => { |
| 227 | let func = self.udf_gen.gen_unfold(unfold)?; |
| 228 | stream = stream.flat_map_with_name("Unfold", move |input| func.exec(input))?; |
| 229 | } |
| 230 | OpKind::Limit(limit) => { |
| 231 | let range = limit.range.ok_or_else(|| { |
| 232 | FnGenError::from(ParsePbError::EmptyFieldError("pb::Limit::range".to_string())) |
| 233 | })?; |
| 234 | // e.g., `limit(10)` would be translate as `Range{lower=0, upper=10}` |
| 235 | if range.upper <= range.lower || range.lower != 0 { |
| 236 | Err(FnGenError::from(ParsePbError::ParseError(format!( |
| 237 | "range {:?} in Limit Operator", |
| 238 | range |
| 239 | ))))?; |
| 240 | } |
| 241 | stream = stream.limit(range.upper as u32)?; |
| 242 | } |
| 243 | OpKind::OrderBy(order) => { |
| 244 | let cmp = self.udf_gen.gen_cmp(order.clone())?; |
| 245 | if let Some(range) = order.limit { |
| 246 | if range.upper <= range.lower || range.lower != 0 { |
| 247 | Err(FnGenError::from(ParsePbError::ParseError(format!( |
| 248 | "range {:?} in Order Operator", |
| 249 | range |
| 250 | ))))?; |
| 251 | } |
| 252 | stream = stream.sort_limit_by(range.upper as u32, move |a, b| cmp.compare(a, b))?; |
| 253 | } else { |
| 254 | stream = stream.sort_by(move |a, b| cmp.compare(a, b))?; |
no test coverage detected