View Javadoc

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      }, "workerProgressThread");
66      writerThread.setDaemon(true);
67      writerThread.start();
68    }
69  
70    /**
71     * Update worker progress and send it
72     */
73    private void updateAndSendProgress() {
74      WorkerProgress.get().updateMemory();
75      jobProgressTracker.updateProgress(WorkerProgress.get());
76    }
77  
78    /**
79     * Stop the thread which writes worker's progress
80     */
81    public void stop() throws InterruptedException {
82      finished = true;
83      writerThread.interrupt();
84      writerThread.join();
85      updateAndSendProgress();
86    }
87  }