
flink任务在yarn失败需要从yarn的归档服务器当中获取运行日志用来排查异常原因正在运行中的任务日志从当前运行的Node节点中读取方便监控任务运行信息。非运行状态的任务日志获取//构建日志请求对象 ContainerLogsRequest request new ContainerLogsRequest(); request.setAppId(report.getApplicationId()); request.setAppFinished(report.getYarnApplicationState().equals(YarnApplicationState.FINISHED)); SetString logs new HashSet(); request.setLogTypes(logs); request.setBytes(Long.MAX_VALUE); request.setAppOwner(report.getUser()); ListApplicationAttemptReport attempts client.getApplicationAttempts(report.getApplicationId()); request.setContainerId(attempts.getLast().getAMContainerId().toString()); //构建log聚合日志工厂类 LogAggregationFileControllerFactory fileControllerFactory new LogAggregationFileControllerFactory(configuration); LogAggregationFileController fileController fileControllerFactory.getFileControllerForRead(report.getApplicationId(), report.getUser()); //获取日志信息 ListContainerLogMeta logMetas fileController.readAggregatedLogsMeta(request); //打印容器日志信息,并且将容器所有内容放入其中 logMetas.forEach(containerLogMeta - { System.out.println(String.format(容器id%s 节点id%s, containerLogMeta.getContainerId(), containerLogMeta.getNodeId())); containerLogMeta.getContainerLogMeta().forEach(fileInfo - { System.out.println(String.format(文件名称%s文件大小%s bytes, fileInfo.getFileName(), fileInfo.getFileSize())); logs.add(fileInfo.getFileName()); }); }); //读取日志 fileController.readAggregatedLogs(request, System.out);运行中状态的任务日志获取YarnConfiguration configuration YarnClientUtil.buildConfig(); YarnClientImpl client YarnClientUtil.client(configuration); //获取yarn上app的状态 ApplicationReport report client.getApplicationReport(ApplicationId.fromString(appId)); OkHttpClient okHttpClient new OkHttpClient(); okHttpClient.setConnectTimeout(3, TimeUnit.SECONDS); okHttpClient.setReadTimeout(15, TimeUnit.SECONDS); //获取当前任务的集群状态 ListApplicationAttemptReport attempts client.getApplicationAttempts(report.getApplicationId()); //遍历集群获取集群下当前容器的报告 for (ApplicationAttemptReport attempt : attempts) { ListContainerReport containers client.getContainers(attempt.getApplicationAttemptId()); //获取当前容器下的日志列表 for (ContainerReport containerReport : containers) { //构造获取容器下的日志文件请求 Request logReq new Request.Builder().url(String.join(/, new String[]{containerReport.getNodeHttpAddress(), ws, v1, node, containers , containerReport.getContainerId().toString(), logs})) .get().header(Accept, application/json).build(); Call call okHttpClient.newCall(logReq); Response response call.execute(); MapString,ListContainerLogFileInfo logType new HashMap(); if (response.isSuccessful()) { String resString response.body().string(); //反序列化日志对象 JSONObject obj JSONObject.parseObject(resString); JSONObject info obj.getJSONObject(containerLogsInfo); if (Objects.nonNull(info)) { String containerId info.getString(containerId); ListContainerLogFileInfo logFileInfos info.getObject(containerLogInfo,new TypeReferenceListContainerLogFileInfo(){}); //System.out.println(String.format(containerLog %s on %s,info.getContainerId(),info.getNodeId())); logType.put(containerId,logFileInfos); //遍历打印容器日志文件信息 logFileInfos.forEach(log - { System.out.println(String.format(容器id%s 文件名%s 大小%s (bytes), containerId, log.getFileName(), log.getFileSize())); }); } } else { throw new RuntimeException(获取日志信息出错); } //遍历日志文件信息并且日志内容 logType.entrySet().forEach(entry-{ String containerId entry.getKey(); ListContainerLogFileInfo fileInfoList entry.getValue(); fileInfoList.forEach(fileInfo - { Request logContentReq new Request.Builder().url(String.join(/, new String[]{containerReport.getNodeHttpAddress(), ws, v1, node, containers , containerId, logs,fileInfo.getFileName()})).build(); try { Response fileResp okHttpClient.newCall(logContentReq).execute(); if(fileResp.isSuccessful()){ IOUtils.copy(fileResp.body().byteStream(), System.out); } } catch (IOException e) { throw new RuntimeException(e); } }); }); } }