dae9963d29ab226e5a0f0c15a09f4b169f0c177b
[giraph.git] / giraph-core / src / main / java / org / apache / giraph / worker / WorkerProgressWriter.java
1 /*
2 * Licensed to the Apache Software Foundation (ASF) under one
3 * or more contributor license agreements. See the NOTICE file
4 * distributed with this work for additional information
5 * regarding copyright ownership. The ASF licenses this file
6 * to you under the Apache License, Version 2.0 (the
7 * "License"); you may not use this file except in compliance
8 * with the License. You may obtain a copy of the License at
9 *
10 * http://www.apache.org/licenses/LICENSE-2.0
11 *
12 * Unless required by applicable law or agreed to in writing, software
13 * distributed under the License is distributed on an "AS IS" BASIS,
14 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15 * See the License for the specific language governing permissions and
16 * limitations under the License.
17 */
18
19 package org.apache.giraph.worker;
20
21 import org.apache.giraph.job.JobProgressTracker;
22 import org.apache.log4j.Logger;
23
24 /**
25 * Class which periodically writes worker's progress to zookeeper
26 */
27 public class WorkerProgressWriter {
28 /** Class logger */
29 private static final Logger LOG =
30 Logger.getLogger(WorkerProgressWriter.class);
31 /** How often to update worker's progress */
32 private static final int WRITE_UPDATE_PERIOD_MILLISECONDS = 10 * 1000;
33
34 /** Job progress tracker */
35 private final JobProgressTracker jobProgressTracker;
36 /** Thread which writes worker's progress */
37 private final Thread writerThread;
38 /** Whether worker finished application */
39 private volatile boolean finished = false;
40
41 /**
42 * Constructor, starts separate thread to periodically update worker's
43 * progress
44 *
45 * @param jobProgressTracker JobProgressTracker to report job progress to
46 */
47 public WorkerProgressWriter(JobProgressTracker jobProgressTracker) {
48 this.jobProgressTracker = jobProgressTracker;
49 writerThread = new Thread(new Runnable() {
50 @Override
51 public void run() {
52 try {
53 while (!finished) {
54 updateAndSendProgress();
55 double factor = 1 + Math.random();
56 Thread.sleep((long) (WRITE_UPDATE_PERIOD_MILLISECONDS * factor));
57 }
58 } catch (InterruptedException e) {
59 // Thread is interrupted when stop is called, we can just log this
60 if (LOG.isInfoEnabled()) {
61 LOG.info("run: WorkerProgressWriter interrupted");
62 }
63 }
64 }
65 });
66 writerThread.start();
67 }
68
69 /**
70 * Update worker progress and send it
71 */
72 private void updateAndSendProgress() {
73 WorkerProgress.get().updateMemory();
74 jobProgressTracker.updateProgress(WorkerProgress.get());
75 }
76
77 /**
78 * Stop the thread which writes worker's progress
79 */
80 public void stop() throws InterruptedException {
81 finished = true;
82 writerThread.interrupt();
83 writerThread.join();
84 updateAndSendProgress();
85 }
86 }