[
https://issues.apache.org/jira/browse/TEZ-4018?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17185577#comment-17185577
]
Nemon Lou commented on TEZ-4018:
--------------------------------
{{Hi [~yingdachen] ,conditional vertex feature sounds useful for join strategy
selection.}}
{{Spark achieve adaptive join selection through a different approach in
SPARK-9851.Via a submitMapStages API.Which need more work for client to
}}{{submit}}{{ stages and schedule between them.}}
{{Any update on this?Thanks.}}
> Allow conditional vertex in DAG execution
> -----------------------------------------
>
> Key: TEZ-4018
> URL: https://issues.apache.org/jira/browse/TEZ-4018
> Project: Apache Tez
> Issue Type: New Feature
> Reporter: Yingda Chen
> Assignee: Yingda Chen
> Priority: Major
>
> A high-level description is provided here for now, *will follow up with a
> proper design doc later.* We welcome any comments/suggests/thoughts from the
> community about this feature.
> We have encountered a few application scenarios for dynamic (logical) DAG
> here in our system. A typical one is for a distributed query to be able to
> dynamically choose among two execution paths, say among hash join and sorted
> merge join.
> This can be solved by allowing TEZ to execute "conditional DAG". By that we
> mean that a DAG may have some conditional vertices: several conditional
> vertices can form a conditional group, and inside each group, only one will
> be chosen for execution at runtime.
> To allow decision at runtime, each conditional group will be associated with
> a “control vertex”. A control vertex can be a pure virtual component that
> lives only on AM with its VertexImpl and VetexManager, but has no associated
> tasks (DoP = 0). It can be also be extended to have physical tasks associated
> with it, in the case where intensive computation may be required to make a
> control decision.
> A upstream vertex (that produces input data to vertices in downstream
> conditional group) will be connected to the control vertex, as well as all
> conditional vertices in the group at the same time. This allows its VMEs and
> DMEs to be sent to all of them. Upon receipt of (enough) VMEs, the control
> vertex would be able to gather enough runtime statistics and determine (by
> user-supplied logic) which downstream vertex should be scheduled (and which
> should be skipped). Such decision will effectively “uncondition” the DAG and
> determines the sub-graph that is actually being executed.
> Such conditional vertices can be useful to enable scenarios such as
> conditional join, where a query can choose between hash join and sorted merge
> join at runtime, base on the precise runtime statistics (e.g., output size)
> of upstream mapper.
--
This message was sent by Atlassian Jira
(v8.3.4#803005)