/
StreamingPlanSection.java
78 lines (69 loc) · 3.07 KB
/
StreamingPlanSection.java
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
/*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.facebook.presto.execution.scheduler;
import com.facebook.presto.sql.planner.SubPlan;
import com.facebook.presto.sql.planner.plan.PlanFragmentId;
import com.facebook.presto.sql.planner.plan.RemoteSourceNode;
import com.google.common.collect.ImmutableList;
import java.util.List;
import java.util.Set;
import static com.google.common.collect.ImmutableList.toImmutableList;
import static com.google.common.collect.ImmutableSet.toImmutableSet;
import static java.util.Objects.requireNonNull;
public class StreamingPlanSection
{
private final StreamingSubPlan plan;
// materialized exchange children
private final List<StreamingPlanSection> children;
public StreamingPlanSection(StreamingSubPlan plan, List<StreamingPlanSection> children)
{
this.plan = requireNonNull(plan, "plan is null");
this.children = ImmutableList.copyOf(requireNonNull(children, "children is null"));
}
public StreamingSubPlan getPlan()
{
return plan;
}
public List<StreamingPlanSection> getChildren()
{
return children;
}
public static StreamingPlanSection extractStreamingSections(SubPlan subPlan)
{
ImmutableList.Builder<SubPlan> materializedExchangeChildren = ImmutableList.builder();
StreamingSubPlan streamingSection = extractStreamingSection(subPlan, materializedExchangeChildren);
return new StreamingPlanSection(
streamingSection,
materializedExchangeChildren.build().stream()
.map(StreamingPlanSection::extractStreamingSections)
.collect(toImmutableList()));
}
private static StreamingSubPlan extractStreamingSection(SubPlan subPlan, ImmutableList.Builder<SubPlan> materializedExchangeChildren)
{
ImmutableList.Builder<StreamingSubPlan> streamingSources = ImmutableList.builder();
Set<PlanFragmentId> streamingFragmentIds = subPlan.getFragment().getRemoteSourceNodes().stream()
.map(RemoteSourceNode::getSourceFragmentIds)
.flatMap(List::stream)
.collect(toImmutableSet());
for (SubPlan child : subPlan.getChildren()) {
if (streamingFragmentIds.contains(child.getFragment().getId())) {
streamingSources.add(extractStreamingSection(child, materializedExchangeChildren));
}
else {
materializedExchangeChildren.add(child);
}
}
return new StreamingSubPlan(subPlan.getFragment(), streamingSources.build());
}
}