1 | |
|
2 | |
|
3 | |
|
4 | |
|
5 | |
|
6 | |
|
7 | |
|
8 | |
|
9 | |
|
10 | |
|
11 | |
|
12 | |
|
13 | |
|
14 | |
|
15 | |
|
16 | |
|
17 | |
|
18 | |
|
19 | |
package org.apache.giraph.comm.netty; |
20 | |
|
21 | |
import org.apache.giraph.bsp.CentralizedServiceWorker; |
22 | |
import org.apache.giraph.comm.WorkerClient; |
23 | |
import org.apache.giraph.comm.flow_control.FlowControl; |
24 | |
import org.apache.giraph.comm.requests.RequestType; |
25 | |
import org.apache.giraph.comm.requests.WritableRequest; |
26 | |
import org.apache.giraph.conf.ImmutableClassesGiraphConfiguration; |
27 | |
import org.apache.giraph.graph.TaskInfo; |
28 | |
import org.apache.giraph.metrics.GiraphMetrics; |
29 | |
import org.apache.giraph.metrics.MetricNames; |
30 | |
import org.apache.giraph.metrics.ResetSuperstepMetricsObserver; |
31 | |
import org.apache.giraph.metrics.SuperstepMetricsRegistry; |
32 | |
import org.apache.giraph.partition.PartitionOwner; |
33 | |
import org.apache.giraph.worker.WorkerInfo; |
34 | |
import org.apache.hadoop.io.Writable; |
35 | |
import org.apache.hadoop.io.WritableComparable; |
36 | |
import org.apache.hadoop.mapreduce.Mapper; |
37 | |
import org.apache.log4j.Logger; |
38 | |
|
39 | |
import com.google.common.collect.Lists; |
40 | |
import com.google.common.collect.Maps; |
41 | |
import com.yammer.metrics.core.Counter; |
42 | |
|
43 | |
import java.io.IOException; |
44 | |
import java.util.List; |
45 | |
import java.util.Map; |
46 | |
|
47 | |
|
48 | |
|
49 | |
|
50 | |
|
51 | |
|
52 | |
|
53 | |
|
54 | |
|
55 | |
@SuppressWarnings("rawtypes") |
56 | |
public class NettyWorkerClient<I extends WritableComparable, |
57 | |
V extends Writable, E extends Writable> implements |
58 | |
WorkerClient<I, V, E>, ResetSuperstepMetricsObserver { |
59 | |
|
60 | 0 | private static final Logger LOG = Logger.getLogger(NettyWorkerClient.class); |
61 | |
|
62 | |
private final ImmutableClassesGiraphConfiguration<I, V, E> conf; |
63 | |
|
64 | |
private final NettyClient nettyClient; |
65 | |
|
66 | |
private final CentralizedServiceWorker<I, V, E> service; |
67 | |
|
68 | |
|
69 | |
|
70 | |
private final Map<RequestType, Counter> superstepRequestCounters; |
71 | |
|
72 | |
|
73 | |
|
74 | |
|
75 | |
|
76 | |
|
77 | |
|
78 | |
|
79 | |
|
80 | |
|
81 | |
public NettyWorkerClient( |
82 | |
Mapper<?, ?, ?, ?>.Context context, |
83 | |
ImmutableClassesGiraphConfiguration<I, V, E> configuration, |
84 | |
CentralizedServiceWorker<I, V, E> service, |
85 | 0 | Thread.UncaughtExceptionHandler exceptionHandler) { |
86 | 0 | this.nettyClient = |
87 | 0 | new NettyClient(context, configuration, service.getWorkerInfo(), |
88 | |
exceptionHandler); |
89 | 0 | this.conf = configuration; |
90 | 0 | this.service = service; |
91 | 0 | this.superstepRequestCounters = Maps.newHashMap(); |
92 | 0 | GiraphMetrics.get().addSuperstepResetObserver(this); |
93 | 0 | } |
94 | |
|
95 | |
@Override |
96 | |
public void newSuperstep(SuperstepMetricsRegistry metrics) { |
97 | 0 | superstepRequestCounters.clear(); |
98 | 0 | superstepRequestCounters.put(RequestType.SEND_WORKER_VERTICES_REQUEST, |
99 | 0 | metrics.getCounter(MetricNames.SEND_VERTEX_REQUESTS)); |
100 | 0 | superstepRequestCounters.put(RequestType.SEND_WORKER_MESSAGES_REQUEST, |
101 | 0 | metrics.getCounter(MetricNames.SEND_WORKER_MESSAGES_REQUESTS)); |
102 | 0 | superstepRequestCounters.put( |
103 | |
RequestType.SEND_PARTITION_CURRENT_MESSAGES_REQUEST, |
104 | 0 | metrics.getCounter( |
105 | |
MetricNames.SEND_PARTITION_CURRENT_MESSAGES_REQUESTS)); |
106 | 0 | superstepRequestCounters.put(RequestType.SEND_PARTITION_MUTATIONS_REQUEST, |
107 | 0 | metrics.getCounter(MetricNames.SEND_PARTITION_MUTATIONS_REQUESTS)); |
108 | 0 | superstepRequestCounters.put(RequestType.SEND_WORKER_AGGREGATORS_REQUEST, |
109 | 0 | metrics.getCounter(MetricNames.SEND_WORKER_AGGREGATORS_REQUESTS)); |
110 | 0 | superstepRequestCounters.put(RequestType.SEND_AGGREGATORS_TO_MASTER_REQUEST, |
111 | 0 | metrics.getCounter(MetricNames.SEND_AGGREGATORS_TO_MASTER_REQUESTS)); |
112 | 0 | superstepRequestCounters.put(RequestType.SEND_AGGREGATORS_TO_OWNER_REQUEST, |
113 | 0 | metrics.getCounter(MetricNames.SEND_AGGREGATORS_TO_OWNER_REQUESTS)); |
114 | 0 | superstepRequestCounters.put(RequestType.SEND_AGGREGATORS_TO_WORKER_REQUEST, |
115 | 0 | metrics.getCounter(MetricNames.SEND_AGGREGATORS_TO_WORKER_REQUESTS)); |
116 | 0 | } |
117 | |
|
118 | |
public CentralizedServiceWorker<I, V, E> getService() { |
119 | 0 | return service; |
120 | |
} |
121 | |
|
122 | |
@Override |
123 | |
public void openConnections() { |
124 | 0 | List<TaskInfo> addresses = Lists.newArrayListWithCapacity( |
125 | 0 | service.getWorkerInfoList().size()); |
126 | 0 | for (WorkerInfo info : service.getWorkerInfoList()) { |
127 | |
|
128 | 0 | if (service.getWorkerInfo().getTaskId() != info.getTaskId()) { |
129 | 0 | addresses.add(info); |
130 | |
} |
131 | 0 | } |
132 | 0 | addresses.add(service.getMasterInfo()); |
133 | 0 | nettyClient.connectAllAddresses(addresses); |
134 | 0 | } |
135 | |
|
136 | |
@Override |
137 | |
public PartitionOwner getVertexPartitionOwner(I vertexId) { |
138 | 0 | return service.getVertexPartitionOwner(vertexId); |
139 | |
} |
140 | |
|
141 | |
@Override |
142 | |
public void sendWritableRequest(int destTaskId, |
143 | |
WritableRequest request) { |
144 | 0 | Counter counter = superstepRequestCounters.get(request.getType()); |
145 | 0 | if (counter != null) { |
146 | 0 | counter.inc(); |
147 | |
} |
148 | 0 | nettyClient.sendWritableRequest(destTaskId, request); |
149 | 0 | } |
150 | |
|
151 | |
@Override |
152 | |
public void waitAllRequests() { |
153 | 0 | nettyClient.waitAllRequests(); |
154 | 0 | } |
155 | |
|
156 | |
@Override |
157 | |
public void closeConnections() throws IOException { |
158 | 0 | nettyClient.stop(); |
159 | 0 | } |
160 | |
|
161 | |
|
162 | |
|
163 | |
|
164 | |
|
165 | |
|
166 | |
|
167 | |
@Override |
168 | |
public void setup(boolean authenticate) { |
169 | 0 | openConnections(); |
170 | 0 | if (authenticate) { |
171 | 0 | authenticate(); |
172 | |
} |
173 | 0 | } |
174 | |
|
175 | |
|
176 | |
|
177 | |
|
178 | |
@Override |
179 | |
public void authenticate() { |
180 | 0 | nettyClient.authenticate(); |
181 | 0 | } |
182 | |
|
183 | |
|
184 | |
|
185 | |
@Override |
186 | |
public FlowControl getFlowControl() { |
187 | 0 | return nettyClient.getFlowControl(); |
188 | |
} |
189 | |
} |