org.apache.hadoop.mapred.LocalJobRunner#org.apache.hadoop.mapreduce.protocol.ClientProtocol源码实例Demo

下面列出了org.apache.hadoop.mapred.LocalJobRunner#org.apache.hadoop.mapreduce.protocol.ClientProtocol 实例代码,或者点击链接到github查看源代码,也可以在右侧发表评论。

源代码1 项目: hadoop   文件: TestJob.java
@Test
public void testJobToString() throws IOException, InterruptedException {
  Cluster cluster = mock(Cluster.class);
  ClientProtocol client = mock(ClientProtocol.class);
  when(cluster.getClient()).thenReturn(client);
  JobID jobid = new JobID("1014873536921", 6);
  JobStatus status = new JobStatus(jobid, 0.0f, 0.0f, 0.0f, 0.0f,
      State.FAILED, JobPriority.NORMAL, "root", "TestJobToString",
      "job file", "tracking url");
  when(client.getJobStatus(jobid)).thenReturn(status);
  when(client.getTaskReports(jobid, TaskType.MAP)).thenReturn(
      new TaskReport[0]);
  when(client.getTaskReports(jobid, TaskType.REDUCE)).thenReturn(
      new TaskReport[0]);
  when(client.getTaskCompletionEvents(jobid, 0, 10)).thenReturn(
      new TaskCompletionEvent[0]);
  Job job = Job.getInstance(cluster, status, new JobConf());
  Assert.assertNotNull(job.toString());
}
 
源代码2 项目: big-c   文件: TestJob.java
@Test
public void testJobToString() throws IOException, InterruptedException {
  Cluster cluster = mock(Cluster.class);
  ClientProtocol client = mock(ClientProtocol.class);
  when(cluster.getClient()).thenReturn(client);
  JobID jobid = new JobID("1014873536921", 6);
  JobStatus status = new JobStatus(jobid, 0.0f, 0.0f, 0.0f, 0.0f,
      State.FAILED, JobPriority.NORMAL, "root", "TestJobToString",
      "job file", "tracking url");
  when(client.getJobStatus(jobid)).thenReturn(status);
  when(client.getTaskReports(jobid, TaskType.MAP)).thenReturn(
      new TaskReport[0]);
  when(client.getTaskReports(jobid, TaskType.REDUCE)).thenReturn(
      new TaskReport[0]);
  when(client.getTaskCompletionEvents(jobid, 0, 10)).thenReturn(
      new TaskCompletionEvent[0]);
  Job job = Job.getInstance(cluster, status, new JobConf());
  Assert.assertNotNull(job.toString());
}
 
源代码3 项目: hadoop   文件: YarnClientProtocolProvider.java
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  if (MRConfig.YARN_FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME))) {
    return new YARNRunner(conf);
  }
  return null;
}
 
源代码4 项目: hadoop   文件: LocalClientProtocolProvider.java
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  String framework =
      conf.get(MRConfig.FRAMEWORK_NAME, MRConfig.LOCAL_FRAMEWORK_NAME);
  if (!MRConfig.LOCAL_FRAMEWORK_NAME.equals(framework)) {
    return null;
  }
  conf.setInt(JobContext.NUM_MAPS, 1);

  return new LocalJobRunner(conf);
}
 
源代码5 项目: hadoop   文件: TestJobMonitorAndPrint.java
@Before
public void setUp() throws IOException {
  conf = new Configuration();
  clientProtocol = mock(ClientProtocol.class);
  Cluster cluster = mock(Cluster.class);
  when(cluster.getConf()).thenReturn(conf);
  when(cluster.getClient()).thenReturn(clientProtocol);
  JobStatus jobStatus = new JobStatus(new JobID("job_000", 1), 0f, 0f, 0f, 0f, 
      State.RUNNING, JobPriority.HIGH, "tmp-user", "tmp-jobname", 
      "tmp-jobfile", "tmp-url");
  job = Job.getInstance(cluster, jobStatus, conf);
  job = spy(job);
}
 
源代码6 项目: big-c   文件: YarnClientProtocolProvider.java
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  if (MRConfig.YARN_FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME))) {
    return new YARNRunner(conf);
  }
  return null;
}
 
源代码7 项目: big-c   文件: LocalClientProtocolProvider.java
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  String framework =
      conf.get(MRConfig.FRAMEWORK_NAME, MRConfig.LOCAL_FRAMEWORK_NAME);
  if (!MRConfig.LOCAL_FRAMEWORK_NAME.equals(framework)) {
    return null;
  }
  conf.setInt(JobContext.NUM_MAPS, 1);

  return new LocalJobRunner(conf);
}
 
源代码8 项目: big-c   文件: TestJobMonitorAndPrint.java
@Before
public void setUp() throws IOException {
  conf = new Configuration();
  clientProtocol = mock(ClientProtocol.class);
  Cluster cluster = mock(Cluster.class);
  when(cluster.getConf()).thenReturn(conf);
  when(cluster.getClient()).thenReturn(clientProtocol);
  JobStatus jobStatus = new JobStatus(new JobID("job_000", 1), 0f, 0f, 0f, 0f, 
      State.RUNNING, JobPriority.HIGH, "tmp-user", "tmp-jobname", 
      "tmp-jobfile", "tmp-url");
  job = Job.getInstance(cluster, jobStatus, conf);
  job = spy(job);
}
 
/** {@inheritDoc} */
@Override public ClientProtocol create(Configuration conf) throws IOException {
    if (FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME))) {
        Collection<String> addrs = conf.getTrimmedStringCollection(MRConfig.MASTER_ADDRESS);

        if (F.isEmpty(addrs))
            throw new IOException("Failed to create client protocol because Ignite node addresses are not " +
                "specified (did you set " + MRConfig.MASTER_ADDRESS + " property?).");

        if (F.contains(addrs, "local"))
            throw new IOException("Local execution mode is not supported, please point " +
                MRConfig.MASTER_ADDRESS + " to real Ignite nodes.");

        Collection<String> addrs0 = new ArrayList<>(addrs.size());

        // Set up port by default if need
        for (String addr : addrs) {
            if (!addr.contains(":"))
                addrs0.add(addr + ':' + ConnectorConfiguration.DFLT_TCP_PORT);
            else
                addrs0.add(addr);
        }

        return new HadoopClientProtocol(conf, client(conf.get(MRConfig.MASTER_ADDRESS), addrs0));
    }

    return null;
}
 
/** {@inheritDoc} */
@Override public ClientProtocol create(InetSocketAddress addr, Configuration conf) throws IOException {
    if (FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME)))
        return createProtocol(addr.getHostString() + ":" + addr.getPort(), conf);

    return null;
}
 
/** {@inheritDoc} */
@Override public void close(ClientProtocol cliProto) throws IOException {
    if (cliProto instanceof HadoopClientProtocol) {
        MapReduceClient cli = ((HadoopClientProtocol)cliProto).client();

        if (cli.release())
            cliMap.remove(cli.cluster(), cli);
    }
}
 
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  if (MRConfig.YARN_TEZ_FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME))) {
    return new YARNRunner(conf);
  }
  return null;
}
 
源代码13 项目: tez   文件: YarnTezClientProtocolProvider.java
@Override
public ClientProtocol create(Configuration conf) throws IOException {
  if (MRConfig.YARN_TEZ_FRAMEWORK_NAME.equals(conf.get(MRConfig.FRAMEWORK_NAME))) {
    return new YARNRunner(conf);
  }
  return null;
}
 
源代码14 项目: hadoop   文件: YarnClientProtocolProvider.java
@Override
public ClientProtocol create(InetSocketAddress addr, Configuration conf)
    throws IOException {
  return create(conf);
}
 
源代码15 项目: hadoop   文件: YarnClientProtocolProvider.java
@Override
public void close(ClientProtocol clientProtocol) throws IOException {
  // nothing to do
}
 
源代码16 项目: hadoop   文件: LocalClientProtocolProvider.java
@Override
public ClientProtocol create(InetSocketAddress addr, Configuration conf) {
  return null; // LocalJobRunner doesn't use a socket
}
 
源代码17 项目: hadoop   文件: LocalClientProtocolProvider.java
@Override
public void close(ClientProtocol clientProtocol) {
  // no clean up required
}
 
源代码18 项目: hadoop   文件: LocalJobRunner.java
public long getProtocolVersion(String protocol, long clientVersion) {
  return ClientProtocol.versionID;
}
 
源代码19 项目: hadoop   文件: Job.java
/** Only for mocking via unit tests. */
@Private
public JobSubmitter getJobSubmitter(FileSystem fs, 
    ClientProtocol submitClient) throws IOException {
  return new JobSubmitter(fs, submitClient);
}
 
源代码20 项目: hadoop   文件: JobSubmitter.java
JobSubmitter(FileSystem submitFs, ClientProtocol submitClient) 
throws IOException {
  this.submitClient = submitClient;
  this.jtFs = submitFs;
}
 
源代码21 项目: hadoop   文件: Cluster.java
private void initialize(InetSocketAddress jobTrackAddr, Configuration conf)
    throws IOException {

  synchronized (frameworkLoader) {
    for (ClientProtocolProvider provider : frameworkLoader) {
      LOG.debug("Trying ClientProtocolProvider : "
          + provider.getClass().getName());
      ClientProtocol clientProtocol = null; 
      try {
        if (jobTrackAddr == null) {
          clientProtocol = provider.create(conf);
        } else {
          clientProtocol = provider.create(jobTrackAddr, conf);
        }

        if (clientProtocol != null) {
          clientProtocolProvider = provider;
          client = clientProtocol;
          LOG.debug("Picked " + provider.getClass().getName()
              + " as the ClientProtocolProvider");
          break;
        }
        else {
          LOG.debug("Cannot pick " + provider.getClass().getName()
              + " as the ClientProtocolProvider - returned null protocol");
        }
      } 
      catch (Exception e) {
        LOG.info("Failed to use " + provider.getClass().getName()
            + " due to error: " + e.getMessage());
      }
    }
  }

  if (null == clientProtocolProvider || null == client) {
    throw new IOException(
        "Cannot initialize Cluster. Please check your configuration for "
            + MRConfig.FRAMEWORK_NAME
            + " and the correspond server addresses.");
  }
}
 
源代码22 项目: hadoop   文件: Cluster.java
ClientProtocol getClient() {
  return client;
}
 
源代码23 项目: big-c   文件: YarnClientProtocolProvider.java
@Override
public ClientProtocol create(InetSocketAddress addr, Configuration conf)
    throws IOException {
  return create(conf);
}
 
源代码24 项目: big-c   文件: YarnClientProtocolProvider.java
@Override
public void close(ClientProtocol clientProtocol) throws IOException {
  // nothing to do
}
 
源代码25 项目: big-c   文件: LocalClientProtocolProvider.java
@Override
public ClientProtocol create(InetSocketAddress addr, Configuration conf) {
  return null; // LocalJobRunner doesn't use a socket
}
 
源代码26 项目: big-c   文件: LocalClientProtocolProvider.java
@Override
public void close(ClientProtocol clientProtocol) {
  // no clean up required
}
 
源代码27 项目: big-c   文件: LocalJobRunner.java
public long getProtocolVersion(String protocol, long clientVersion) {
  return ClientProtocol.versionID;
}
 
源代码28 项目: big-c   文件: Job.java
/** Only for mocking via unit tests. */
@Private
public JobSubmitter getJobSubmitter(FileSystem fs, 
    ClientProtocol submitClient) throws IOException {
  return new JobSubmitter(fs, submitClient);
}
 
源代码29 项目: big-c   文件: JobSubmitter.java
JobSubmitter(FileSystem submitFs, ClientProtocol submitClient) 
throws IOException {
  this.submitClient = submitClient;
  this.jtFs = submitFs;
}
 
源代码30 项目: big-c   文件: Cluster.java
private void initialize(InetSocketAddress jobTrackAddr, Configuration conf)
    throws IOException {

  synchronized (frameworkLoader) {
    for (ClientProtocolProvider provider : frameworkLoader) {
      LOG.debug("Trying ClientProtocolProvider : "
          + provider.getClass().getName());
      ClientProtocol clientProtocol = null; 
      try {
        if (jobTrackAddr == null) {
          clientProtocol = provider.create(conf);
        } else {
          clientProtocol = provider.create(jobTrackAddr, conf);
        }

        if (clientProtocol != null) {
          clientProtocolProvider = provider;
          client = clientProtocol;
          LOG.debug("Picked " + provider.getClass().getName()
              + " as the ClientProtocolProvider");
          break;
        }
        else {
          LOG.debug("Cannot pick " + provider.getClass().getName()
              + " as the ClientProtocolProvider - returned null protocol");
        }
      } 
      catch (Exception e) {
        LOG.info("Failed to use " + provider.getClass().getName()
            + " due to error: " + e.getMessage());
      }
    }
  }

  if (null == clientProtocolProvider || null == client) {
    throw new IOException(
        "Cannot initialize Cluster. Please check your configuration for "
            + MRConfig.FRAMEWORK_NAME
            + " and the correspond server addresses.");
  }
}