Flink使用RestApi

flink是一个非常好用的流任务计算框架, 这次我们来试用flink的restApi来提交任务. 主要阐述几个常用的restapi, 包括上传jar包, 查询jar包, 提交任务, 查询任务, 删除任务等,
其它的比如删除jar包, 查询jobmanager, 查询taskmanager等等, 类推就可以得出了, 不在这里进行重复介绍了

1, 上传jar包

1 public static boolean uploadJar( File jarFile) { 2 RequestBody requestBody = new MultipartBody.Builder() 3 .setType(MultipartBody.FORM) 4 .addFormDataPart("file", jarFile.getName(), 5 RequestBody.create(MediaType.parse("multipart/form-data"), jarFile)) 6 .build(); 7 8 Request request = new Request.Builder() 9 .url("http://host:port/jars/upload") 10 .addHeader(userAgent, userAgentVal) 11 .post(requestBody) 12 .build(); 13 Response resp = OkHttpUtils.execute(request); 14 if (OK == resp.code()) { 15 JSONObject body = JSON.parseObject(resp.body().string()); 16 if ("success".equals(body.getString("status"))) { 17 return true; 18 } 19 } 20 return false; 21 }

2, 查询jar包

1 Request request = new Request.Builder() 2 .url("http://host:port/jars") 3 .addHeader(userAgent, userAgentVal) 4 .get() 5 .build(); 6 Response response = OkHttpUtils.execute(request); 7 String body = response.body().string();

3,提交任务
(特别提示: 提交任务时, Main方法中,容易出现参数解析异常, 为了解决这一个问题, 强烈建议, 对参数进行编码转换, 对programArgs参数进行URLEncoder.encode(参数值, “utf-8”), 然后再在flink运行jar包, 进行解码.

1 String baseUrl = "http://host:port/jars/${jarId}/run"; 2 Map<String, String> params = new HashMap<>(); 3 params.put("programArgs", "xxxxxx"); 4 params.put("entryClass", "com.xx.oo.JsonMain"); 5 params.put("parallelism", "2"); 6 params.put("savepointPath", null); 7 Request request = new Request.Builder() 8 .url(baseUrl) 9 .addHeader(userAgent, userAgentVal) 10 .post(RequestBody.create(JSON.toJSONString(params), MEDIA_TYPE_JSON)) 11 .build(); 12 Response resp = OkHttpUtils.execute(request); 13 String respBody = resp.body().string(); 14 if (OK == resp.code()) { 15 JSONObject body = JSON.parseObject(respBody); 16 return body.getString("jobid"); 17 }

4,查询任务

1 String url= "http://host:port/jobs"; 2 Request request = new Request.Builder() 3 .url(url) 4 .addHeader(userAgent, userAgentVal) 5 .get() 6 .build(); 7 Response resp = OkHttpUtils.execute(request); 8 if (OK == resp.code()) { 9 JSONObject body = JSON.parseObject(resp.body().string()); 10 if (body.containsKey("jobs")) { 11 JSONArray jobs = body.getJSONArray("jobs"); 12 for (int i = 0; i < jobs.size(); i++) { 13 JSONObject jsb = jobs.getJSONObject(i); 14 String id = jsb.getString("id"); 15 String status = jsb.getString("status"); 16 } 17 } 18 }else{ 19 logger.error("queryJobByHttp "+resp.body().string()); 20 }

4,删除任务

1 String url= "http://host:port/jobs/${jobId}"; 2 Request request = new Request.Builder() 3 .url(baseUrl) 4 .addHeader(userAgent, userAgentVal) 5 .patch(RequestBody.create("{}", MEDIA_TYPE_JSON)) 6 .build(); 7 Response resp = OkHttpUtils.execute(request); 8 if (ACCEPTED == resp.code()) { 9 return jobId; 10 } 11 return null;
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang