Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- /**
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you 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 org.apache.hama.bsp;
- import java.io.IOException;
- import java.net.InetAddress;
- import java.net.InetSocketAddress;
- import java.util.Iterator;
- import java.util.List;
- import java.util.Map.Entry;
- import junit.framework.TestCase;
- import org.apache.commons.logging.Log;
- import org.apache.commons.logging.LogFactory;
- import org.apache.hadoop.conf.Configuration;
- import org.apache.hadoop.fs.FSDataInputStream;
- import org.apache.hadoop.fs.FileSystem;
- import org.apache.hadoop.fs.Path;
- import org.apache.hadoop.io.NullWritable;
- import org.apache.hadoop.io.Text;
- import org.apache.hadoop.io.Writable;
- import org.apache.hadoop.ipc.RPC;
- import org.apache.hadoop.ipc.Server;
- import org.apache.hama.Constants;
- import org.apache.hama.HamaConfiguration;
- import org.apache.hama.bsp.Counters.Counter;
- import org.apache.hama.bsp.TestBSPTaskFaults.MinimalGroomServer;
- import org.apache.hama.bsp.ft.CheckpointService;
- import org.apache.hama.bsp.ft.IFaultTolerantPeerService;
- import org.apache.hama.bsp.message.HadoopMessageManager;
- import org.apache.hama.bsp.message.MessageManager;
- import org.apache.hama.bsp.message.MessageManagerFactory;
- import org.apache.hama.bsp.message.MessageQueue;
- import org.apache.hama.bsp.message.type.ByteMessage;
- import org.apache.hama.bsp.sync.BSPPeerSyncClient;
- import org.apache.hama.bsp.sync.PeerSyncClient;
- import org.apache.hama.bsp.sync.SyncClient;
- import org.apache.hama.bsp.sync.SyncEvent;
- import org.apache.hama.bsp.sync.SyncEventListener;
- import org.apache.hama.bsp.sync.SyncException;
- import org.apache.hama.bsp.sync.SyncServiceFactory;
- import org.apache.hama.bsp.sync.ZooKeeperSyncClientImpl;
- import org.apache.hama.ipc.BSPPeerProtocol;
- import org.apache.hama.ipc.HamaRPCProtocolVersion;
- import org.apache.hama.util.BSPNetUtils;
- import org.apache.hama.util.KeyValuePair;
- public class TestCheckpoint extends TestCase {
- public static final Log LOG = LogFactory.getLog(TestCheckpoint.class);
- static final String checkpointedDir = "checkpoint/job_201110302255_0001/0/";
- public static class TestMessageManager<Text> implements MessageManager<Writable>{
- List<Text> messageQueue;
- @Override
- public void init(TaskAttemptID attemptId,
- BSPPeer<?, ?, ?, ?, Writable> peer, Configuration conf,
- InetSocketAddress peerAddress) {
- // TODO Auto-generated method stub
- }
- @Override
- public void close() {
- // TODO Auto-generated method stub
- }
- @Override
- public Writable getCurrentMessage() throws IOException {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public void send(String peerName, Writable msg) throws IOException {
- // TODO Auto-generated method stub
- }
- @Override
- public void finishSendPhase() throws IOException {
- // TODO Auto-generated method stub
- }
- @Override
- public Iterator<Entry<InetSocketAddress, MessageQueue<Writable>>> getMessageIterator() {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public void transfer(InetSocketAddress addr,
- BSPMessageBundle<Writable> bundle) throws IOException {
- // TODO Auto-generated method stub
- }
- @Override
- public void clearOutgoingQueues() {
- // TODO Auto-generated method stub
- }
- @Override
- public int getNumCurrentMessages() {
- // TODO Auto-generated method stub
- return 0;
- }
- @Override
- public void loopBackMessage(BSPMessageBundle<? extends Writable> bundle) {
- // TODO Auto-generated method stub
- }
- @Override
- public void replayMessages() {
- // TODO Auto-generated method stub
- }
- }
- public static class TestBSPPeer
- implements BSPPeer<NullWritable, NullWritable, NullWritable, NullWritable, Writable> {
- Configuration conf;
- long superstepCount;
- public TestBSPPeer(Configuration conf){
- this.conf = conf;
- superstepCount = 0;
- }
- @Override
- public void send(String peerName, Writable msg) throws IOException {}
- @Override
- public Writable getCurrentMessage() throws IOException {
- return new Text("data");
- }
- @Override
- public int getNumCurrentMessages() {
- return 1;
- }
- @Override
- public void sync() throws IOException, SyncException, InterruptedException {
- ++superstepCount;
- }
- @Override
- public long getSuperstepCount() {
- // TODO Auto-generated method stub
- return 0;
- }
- @Override
- public String getPeerName() {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public String getPeerName(int index) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public int getPeerIndex() {
- // TODO Auto-generated method stub
- return 1;
- }
- @Override
- public String[] getAllPeerNames() {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public int getNumPeers() {
- // TODO Auto-generated method stub
- return 0;
- }
- @Override
- public void clear() {
- // TODO Auto-generated method stub
- }
- @Override
- public void write(NullWritable key, NullWritable value) throws IOException {
- // TODO Auto-generated method stub
- }
- @Override
- public boolean readNext(NullWritable key, NullWritable value)
- throws IOException {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public KeyValuePair<NullWritable, NullWritable> readNext()
- throws IOException {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public void reopenInput() throws IOException {
- // TODO Auto-generated method stub
- }
- @Override
- public Configuration getConfiguration() {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public Counter getCounter(Enum<?> name) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public Counter getCounter(String group, String name) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public void incrementCounter(Enum<?> key, long amount) {
- // TODO Auto-generated method stub
- }
- @Override
- public void incrementCounter(String group, String counter, long amount) {
- // TODO Auto-generated method stub
- }
- }
- public static class TempSyncClient extends BSPPeerSyncClient {
- @Override
- public String constructKey(BSPJobID jobId, String... args) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public boolean storeInformation(String key, Writable value,
- boolean permanent, SyncEventListener listener) {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public Writable getInformation(String key,
- Class<? extends Writable> classType) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public boolean addKey(String key, boolean permanent,
- SyncEventListener listener) {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public boolean hasKey(String key) {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public String[] getChildKeySet(String key, SyncEventListener listener) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public boolean registerListener(String key, SyncEvent event,
- SyncEventListener listener) {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public boolean remove(String key, SyncEventListener listener) {
- // TODO Auto-generated method stub
- return false;
- }
- @Override
- public void init(Configuration conf, BSPJobID jobId, TaskAttemptID taskId)
- throws Exception {
- // TODO Auto-generated method stub
- }
- @Override
- public void enterBarrier(BSPJobID jobId, TaskAttemptID taskId,
- long superstep) throws SyncException {
- LOG.info("Enter barrier called - " + superstep);
- }
- @Override
- public void leaveBarrier(BSPJobID jobId, TaskAttemptID taskId,
- long superstep) throws SyncException {
- LOG.info("Exit barrier called - " + superstep);
- }
- @Override
- public void register(BSPJobID jobId, TaskAttemptID taskId,
- String hostAddress, long port) {
- // TODO Auto-generated method stub
- }
- @Override
- public String[] getAllPeerNames(TaskAttemptID taskId) {
- // TODO Auto-generated method stub
- return null;
- }
- @Override
- public void deregisterFromBarrier(BSPJobID jobId, TaskAttemptID taskId,
- String hostAddress, long port) {
- // TODO Auto-generated method stub
- }
- @Override
- public void stopServer() {
- // TODO Auto-generated method stub
- }
- @Override
- public void close() throws IOException {
- // TODO Auto-generated method stub
- }
- }
- @SuppressWarnings({ "unchecked", "rawtypes" })
- public void testCheckpoint() throws Exception {
- Configuration config = new Configuration();
- config.set(SyncServiceFactory.SYNC_PEER_CLASS,
- TempSyncClient.class.getName());
- config.set(Constants.FAULT_TOLERANCE_CLASS,
- CheckpointService.class.getName());
- int port = BSPNetUtils.getFreePort(12502);
- LOG.info("Got port = " + port);
- config.set(Constants.PEER_HOST, Constants.DEFAULT_PEER_HOST);
- config.setInt(Constants.PEER_PORT, port);
- config.set("bsp.output.dir", "/tmp/hama-test_out");
- FileSystem dfs = FileSystem.get(config);
- BSPJob job = new BSPJob(new BSPJobID("checkpttest", 1), "/tmp");
- TaskAttemptID taskId = new TaskAttemptID(new TaskID(job.getJobID(), 1), 1);
- //BSPPeerImpl bspTask = new BSPPeerImpl(job, config, dfs, taskId);
- BSPPeer bspTask = new TestBSPPeer(config);
- TestMessageManager<Text> messenger = new TestMessageManager<Text>();
- assertNotNull("BSPPeerImpl should not be null.", bspTask);
- if (dfs.mkdirs(new Path("checkpoint"))) {
- if (dfs.mkdirs(new Path("checkpoint/job_201110302255_0001"))) {
- if (dfs.mkdirs(new Path("checkpoint/job_201110302255_0001/0")))
- ;
- }
- }
- LOG.info("Created bsp peer and other parameters");
- // HadoopMessageManager<? extends Writable> messageManager =
- // (HadoopMessageManager<? extends Writable>) RPC.getProxy(
- // HadoopMessageManager.class, HamaRPCProtocolVersion.versionID,
- // new InetSocketAddress(port), config);
- // InetSocketAddress peerAddress = new InetSocketAddress(port);
- // MessageManager<? extends Writable> messageManager = MessageManagerFactory.getMessageManager(config);
- // messageManager.init(taskId, bspTask, config, peerAddress);
- PeerSyncClient syncClient = (TempSyncClient) SyncServiceFactory
- .getPeerSyncClient(config);
- IFaultTolerantPeerService<Text> service = null;
- boolean initialized = false;
- try{
- service = (new CheckpointService<Text>())
- .constructPeerFaultTolerance(job, (BSPPeer)bspTask, syncClient,
- new InetSocketAddress(port), taskId, -1L , config,
- (MessageManager<Text>)messageManager);
- initialized = true;
- }
- catch(Exception e){
- }
- assertTrue(initialized);
- LOG.info("Initialized correctly");
- assertTrue("Make sure directory is created.",
- dfs.exists(new Path(checkpointedDir)));
- byte[] tmpData = "data".getBytes();
- // BSPMessageBundle bundle = new BSPMessageBundle();
- // bundle.addMessage(new ByteMessage("abc".getBytes(), tmpData));
- // assertNotNull("Message bundle can not be null.", bundle);
- // assertNotNull("Configuration should not be null.", config);
- //
- // messageManager.put(bundle);
- bspTask.sync();
- LOG.info("out of sync");
- //Path path = new Path()
- // bspTask.checkpointSentMessages(checkpointedDir + "/attempt_201110302255_0001_000000_0",
- // bundle);
- FSDataInputStream in = dfs.open(new Path(checkpointedDir
- + "/attempt_201110302255_0001_000000_0"));
- BSPMessageBundle bundleRead = new BSPMessageBundle();
- bundleRead.readFields(in);
- in.close();
- ByteMessage byteMsg = (ByteMessage) (bundleRead.getMessages()).get(0);
- String content = new String(byteMsg.getData());
- LOG.info("Saved checkpointed content is " + content);
- assertTrue("Message content should be the same.", "data".equals(content));
- dfs.delete(new Path("checkpoint"), true);
- }
- // public void testCheckpointInterval() throws Exception {
- //
- // Configuration conf = new Configuration();
- // conf.set("bsp.output.dir", "/tmp/hama-test_out");
- // conf.setClass(SyncServiceFactory.SYNC_PEER_CLASS,
- // LocalBSPRunner.LocalSyncClient.class, SyncClient.class);
- //
- // conf.setBoolean(Constants.CHECKPOINT_ENABLED, false);
- //
- // int port = BSPNetUtils.getFreePort(5000);
- // InetSocketAddress inetAddress = new InetSocketAddress(port);
- // MinimalGroomServer groom = new MinimalGroomServer(conf);
- // Server workerServer = RPC.getServer(groom, inetAddress.getHostName(),
- // inetAddress.getPort(), conf);
- // workerServer.start();
- //
- // LOG.info("Started RPC server");
- // conf.setInt("bsp.groom.rpc.port", inetAddress.getPort());
- // conf.setInt("bsp.peers.num", 1);
- //
- // BSPPeerProtocol umbilical = (BSPPeerProtocol) RPC.getProxy(
- // BSPPeerProtocol.class, HamaRPCProtocolVersion.versionID, inetAddress,
- // conf);
- // LOG.info("Started the proxy connections");
- //
- // TaskAttemptID tid = new TaskAttemptID(new TaskID(new BSPJobID(
- // "job_201110102255", 1), 1), 1);
- //
- // try {
- // BSPJob job = new BSPJob(new HamaConfiguration(conf));
- // job.setOutputPath(TestBSPMasterGroomServer.OUTPUT_PATH);
- // job.setOutputFormat(TextOutputFormat.class);
- // final BSPPeerProtocol proto = (BSPPeerProtocol) RPC.getProxy(
- // BSPPeerProtocol.class, HamaRPCProtocolVersion.versionID,
- // new InetSocketAddress("127.0.0.1", port), conf);
- //
- // BSPTask task = new BSPTask();
- // task.setConf(job);
- //
- // @SuppressWarnings("rawtypes")
- // BSPPeerImpl<?, ?, ?, ?, ?> bspPeer = new BSPPeerImpl(job, conf, tid,
- // proto, 0, null, null, new Counters());
- //
- // bspPeer.setCurrentTaskStatus(new TaskStatus(new BSPJobID(), tid, 1.0f,
- // TaskStatus.State.RUNNING, "running", "127.0.0.1",
- // TaskStatus.Phase.STARTING, new Counters()));
- //
- // assertEquals(bspPeer.isReadyToCheckpoint(), false);
- //
- // conf.setBoolean(Constants.CHECKPOINT_ENABLED, true);
- // conf.setInt(Constants.CHECKPOINT_INTERVAL, 3);
- //
- // bspPeer.sync();
- //
- // LOG.info("Is Ready = " + bspPeer.isReadyToCheckpoint() + " at step "
- // + bspPeer.getSuperstepCount());
- // assertEquals(bspPeer.isReadyToCheckpoint(), false);
- // bspPeer.sync();
- // LOG.info("Is Ready = " + bspPeer.isReadyToCheckpoint() + " at step "
- // + bspPeer.getSuperstepCount());
- // assertEquals(bspPeer.isReadyToCheckpoint(), false);
- // bspPeer.sync();
- // LOG.info("Is Ready = " + bspPeer.isReadyToCheckpoint() + " at step "
- // + bspPeer.getSuperstepCount());
- // assertEquals(bspPeer.isReadyToCheckpoint(), true);
- //
- // job.setCheckPointInterval(5);
- // bspPeer.sync();
- // LOG.info("Is Ready = " + bspPeer.isReadyToCheckpoint() + " at step "
- // + bspPeer.getSuperstepCount());
- // assertEquals(bspPeer.isReadyToCheckpoint(), false);
- // bspPeer.sync();
- // LOG.info("Is Ready = " + bspPeer.isReadyToCheckpoint() + " at step "
- // + bspPeer.getSuperstepCount());
- // assertEquals(bspPeer.isReadyToCheckpoint(), false);
- //
- // } catch (Exception e) {
- // LOG.error("Error testing BSPPeer.", e);
- // } finally {
- // umbilical.close();
- // Thread.sleep(2000);
- // workerServer.stop();
- // Thread.sleep(2000);
- // }
- //
- // }
- }
Advertisement
Add Comment
Please, Sign In to add comment