86.通过Livy的RESTful API接口向CDH集群提交作业

86.1 演示环境介绍

  • Livy版本为:0.4
  • CM和CDH版本为:5.13.1
  • 集群未启用Kerberos

86.2 操作演示

  • 将作业运行的jar包上传到HDFS目录
  • 使用Maven创建Livy示例工程
  • pom文件中添加如下依赖
<dependency>
    <groupId>org.apache.httpcomponents</groupId>
    <artifactId>httpclient</artifactId>
    <version>4.5.4</version>
</dependency>

1.编写示例代码

  • HTTP请求的工具类(HttpUtils.java)
package com.cloudera.utils;
import org.apache.http.HttpEntity;
import org.apache.http.HttpResponse;
import org.apache.http.client.methods.HttpDelete;
import org.apache.http.client.methods.HttpGet;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClients;
import org.apache.http.util.EntityUtils;
import java.io.IOException;
import java.util.Map;
/**
 * package: com.cloudera
 * describe: 封装非Kerberos环境的Http请求工具类
 * creat_user: Fayson
 * email: htechinfo@163.com
 * creat_date: 2019/2/12
 * creat_time: 下午12:16
 * 公众号:碧茂科技
 */
public class HttpUtils {
    /**
     * HttpGET请求
     * @param url
     * @param headers
     * @return
     */
    public static String getAccess(String url, Map<String,String> headers) {
        String result = null;
        CloseableHttpClient httpClient = HttpClients.createDefault();
        HttpGet httpGet = new HttpGet(url);
        if(headers != null && headers.size() > 0){
            headers.forEach((K,V)->httpGet.addHeader(K,V));
        }
        try {
            HttpResponse response = httpClient.execute(httpGet);
            HttpEntity entity = response.getEntity();
            result = EntityUtils.toString(entity);
            System.out.println(result);
        } catch (IOException e) {
            e.printStackTrace();
        }
        return result;
    }
    /**
     * HttpDelete请求
     * @param url
     * @param headers
     * @return
     */
    public  static String deleteAccess(String url, Map<String,String> headers) {
        String result = null;
        CloseableHttpClient httpClient = HttpClients.createDefault();
        HttpDelete httpDelete = new HttpDelete(url);
        if(headers != null && headers.size() > 0){
            headers.forEach((K,V)->httpDelete.addHeader(K,V));
        }
        try {
            HttpResponse response = httpClient.execute(httpDelete);
            HttpEntity entity = response.getEntity();
            result = EntityUtils.toString(entity);
            System.out.println(result);
        } catch (IOException e) {
            e.printStackTrace();
        }
        return result;
    }
    /**
     * HttpPost请求
     * @param url
     * @param headers
     * @param data
     * @return
     */
    public static String postAccess(String url, Map<String,String> headers, String data)  {
        String result = null;
        CloseableHttpClient httpClient = HttpClients.createDefault();
        HttpPost post = new HttpPost(url);
        if(headers != null && headers.size() > 0){
            headers.forEach((K,V)->post.addHeader(K,V));
        }
        try {
            StringEntity entity = new StringEntity(data);
            entity.setContentEncoding("UTF-8");
            entity.setContentType("application/json");
            post.setEntity(entity);
            HttpResponse response = httpClient.execute(post);
            HttpEntity resultEntity = response.getEntity();
            result = EntityUtils.toString(resultEntity);
            System.out.println(result);
            return result;
        } catch (Exception e) {
            e.printStackTrace();
        }
        return result;
    }
}
  • Livy RESTful API调用示例代码
package com.cloudera.nokerberos;
import com.cloudera.utils.HttpUtils;
import java.util.HashMap;
/**
 * package: com.cloudera
 * describe: 通过Java代码调用Livy的RESTful API实现向非Kerberos的CDH集群作业提交
 * creat_user: Fayson
 * email: htechinfo@163.com
 * creat_date: 2019/2/11
 * creat_time: 上午10:50
 * 公众号:碧茂科技
 */
public class AppLivy {
    private static String LIVY_HOST = "http://ip-186-31-7-186.fayson.com:8998";
    public static void main(String[] args) {
        HashMap<String, String> headers = new HashMap<>();
        headers.put("Content-Type", "application/json");
        headers.put("Accept", "application/json");
        headers.put("X-Requested-By", "fayson");
        //创建一个交互式会话
//        String kindJson = "{\"kind\": \"spark\", \"proxyUser\":\"fayson\"}";
//        HttpUtils.postAccess(LIVY_HOST + "/sessions", headers, kindJson);
        //执行code
//        String code = "{\"code\":\"sc.parallelize(1 to 2).count()\"}";
//        HttpUtils.postAccess(LIVY_HOST + "/sessions/1/statements", headers, code);
        //删除会话
//        HttpUtils.deleteAccess(LIVY_HOST + "/sessions/2", headers);
        //封装提交Spark作业的JSON数据
        String submitJob = "{\"className\": \"org.apache.spark.examples.SparkPi\",\"executorMemory\": \"1g\",\"args\": [200],\"file\": \"/fayson-yarn/jars/spark-examples-1.6.0-cdh5.13.1-hadoop2.6.0-cdh5.13.1.jar\", \"proxyUser\":\"fayson\"}";
        //向集群提交Spark作业
        HttpUtils.postAccess(LIVY_HOST + "/batches", headers, submitJob);
        //通过提交作业返回的SessionID获取具体作业的执行状态及APPID
        HttpUtils.getAccess(LIVY_HOST + "/batches/3", headers);
    }
}

2.代码运行

  • 运行AppLivy代码,向集群提交Spark作业,响应结果:
{
  "id": 4,
  "state": "starting",
  "appId": null,
  "appInfo": {
    "driverLogUrl": null,
    "sparkUiUrl": null
  },
  "log": ["stdout: ", "\nstderr: ", "\nYARN Diagnostics: "]
}
  • 获取作业运行状态,将上一步获取到的id传入到如下请求


  • 响应结果:
    • 通过如上返回的结果,我们可以看到作业的APPID
{
  "id": 4,
  "state": "success",
  "appId": "application_1518401384543_0008",
  "appInfo": {
    "driverLogUrl": null,
    "sparkUiUrl": "http://ip-186-31-6-148.fayson.com:8088/proxy/application_1518401384543_0008/"
  },
  "log": ["stdout: ", "WARNING: User-defined SPARK_HOME (/opt/cloudera/parcels/CDH-5.13.1-1.cdh5.13.1.p0.2/lib/spark) overrides detected (/opt/cloudera/parcels/CDH/lib/spark).", "WARNING: Running spark-class from user-defined location.", "\nstderr: ", "\nYARN Diagnostics: "]
}
  • 查看Livy界面提交作业的状态
    • 通过CM和Yarn的8088界面查看作业执行结果

大数据视频推荐:
腾讯课堂
CSDN
大数据语音推荐:
企业级大数据技术应用
大数据机器学习案例之推荐系统
自然语言处理
大数据基础
人工智能:深度学习入门到精通

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容