使用java在aliyun/aws创建E-MapReduce (emr)集群

2024-02-09 15:38

本文主要是介绍使用java在aliyun/aws创建E-MapReduce (emr)集群,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

文章目录

    • 背景
    • 功能点
    • 开撸
      • emr接口
      • emr抽象类
      • 阿里云集群的创建

背景

在上个公司,我的 hera 任务调度系统是运行在本地 cdh 机器上的,并没有使用 aws/aliyun 提供的 emr 服务。所以为了使 hera 能够兼容 emr,就需要使用 java 创建 emr 集群.

功能点

既然要创建集群,肯定也要有等待集群创建完成、销毁集群等操作。

所以功能点大概有

  • 判断集群是否已经创建过
  • 创建集群
  • 等待集群创建完成
  • 集群销毁
  • 获得集群的登录脚本
  • 集群弹性伸缩
  • 等等

开撸

emr接口

为了增加程序的可维护性和扩展性,我提供了一个Emr接口类

/*** desc:** @author scx* @create 2019/04/01*/
public interface Emr {/*** 添加任务接口*/void addJob();/*** 移除任务*/void removeJob();/*** 获得登录命令* @return  登录命令*/String getLogin(String user);}

主要有三个方法,其中 addJobremoveJob 的作用是为了后面进行判断集群是否有任务执行来进行关闭集群的操作。而 getLogin 方法是为了获得集群的 ssh 登录命令

emr抽象类

我希望在抽象类中完成所有的同步操作、集群关闭监听事件、等待集群创建完成 等等公共操作,而子类只关注一些简单实现即可。具体请看源码

/*** desc: emr集群 抽象类** @author scx* @create 2019/04/01*/
public abstract class AbstractEmr implements Emr {/*** 缓存的集群IP*/protected volatile String cacheIp = null;/*** 任务数*/private AtomicInteger taskRunning;/*** 任务计数器*/private AtomicLong taskNum;/*** 缓存的集群Id*/private String cacheClusterId;/*** 关闭集群的调度池*/private ScheduledExecutorService pool;/*** 上次的任务数*/private long cacheTaskNum;/*** 集群是否已经关闭字段*/private volatile boolean clusterTerminate = true;/*** check 集群是否需要关闭返回的future*/private ScheduledFuture<?> clusterWatchFuture;/*** 当前应用的环境 设置不同的参数*/private String env = HeraGlobalEnvironment.getEnv() == null ? "daily" : HeraGlobalEnvironment.getEnv();/*** emr集群的前缀*/protected final String clusterPrefix = "hera-schedule-" + env + "-";/*** 自动扩展冷却时间 单位:秒*/private int coolDown = 300;/*** 扩展百分比*/private int scalePercent = 10;/*** 最少实例数*/private int minCapacity = 1;/*** 最大实例数*/private int maxCapacity = 10;public int getCoolDown() {return coolDown;}public int getScalePercent() {return scalePercent;}public int getMinCapacity() {return minCapacity;}public int getMaxCapacity() {return maxCapacity;}protected String getClusterName() {return clusterPrefix + ActionUtil.getCurrDate();}protected String getEnv() {return env;}/*** 创建集群的方法*/protected void createCluster() {//抽象方法留给子类实现init();//dubbo-check 同步,避免多次创建if (clusterTerminate) {synchronized (this) {if (clusterTerminate) {//启动时可能已经有创建好的集群 首先判断if (notAlive()) {//如果没有创建好的集群,发送创建集群的请求,抽象方法具体留给子类实现cacheClusterId = sendCreateRequest();}//参数初始化clusterTerminate = false;taskRunning = new AtomicInteger(0);taskNum = new AtomicLong(0);//等待集群创建完成waitClusterCompletion();//集群关闭的监视器submitClusterWatch();MonitorLog.info("集群创建完成,可以执行任务了.集群ID为:" + cacheClusterId + ".集群IP为:" + getIp());}}} else {//此时虽然集群被判断为活着,但是也有可能被人为关闭,recheckif (notAlive()) {destroyCluster();createCluster();}}}@Overridepublic void addJob() {//每次调用addJob方法都要执行一下createCluster方法createCluster();//正在执行的任务数+1taskRunning.incrementAndGet();//总任务数计数taskNum.incrementAndGet();}@Overridepublic void removeJob() {//正在执行的任务数-1if (taskRunning != null) {taskRunning.decrementAndGet();}}/*** 判断是否有创建好的集群* @return  结果*/private boolean notAlive() {//getAliveId 抽象方法留给具体子类实现,返回已经创建好的集群idreturn StringUtils.isBlank(cacheClusterId = getAliveId());}/*** 获得emr集群的登录ip/域名* @return  ip/域名*/public String getIp() {//同样dubbo-check 避免多次判断if (cacheIp == null) {synchronized (this) {if (cacheIp == null) {createCluster();// Amazon返回的是域名 ,aliyun目前只返回了master的内网ip,可以自己自定义 getMasterIp 的实现cacheIp = getMasterIp(cacheClusterId);if (StringUtils.isBlank(cacheIp)) {throw new NullPointerException("cacheIp can not be null");}}}}return cacheIp;}/*** 循环检测 ,等待集群创建完成*/private void waitClusterCompletion() {long start = System.currentTimeMillis();long sleepTime = 15 * 1000 * 1000000L;// 必须同步 不能异步while (!isCompletion(cacheClusterId)) {LockSupport.parkNanos(sleepTime);}MonitorLog.info("创建集群:" + cacheClusterId + "耗时:" + (System.currentTimeMillis() - start) + "ms");}/*** createCluster 方法已经同步过* 关于为什么为这里使用的11分钟:我的hera任务调度系统在任务失败后有任务重试,* 重试的间隔时间为10分钟,由于创建集群的时间过长,* 避免任务在重试期间集群被关闭,然后又重新创建,* 所以选择了略大于10的值11。具体时间,大家可以自己衡量*/private void submitClusterWatch() {if (clusterWatchFuture == null) {if (pool == null) {pool = new ScheduledThreadPoolExecutor(1, new NamedThreadFactory("cluster-destroy-watch", true));}//11分钟前的任务数cacheTaskNum = taskNum.get();clusterWatchFuture = pool.scheduleWithFixedDelay(() -> {MonitorLog.info("正在emr集群运行的任务个数:{},十分钟前运行的总任务个数:{},现在运行的总任务个数:{}", taskRunning.get(), cacheTaskNum, taskNum.get());//如果正在执行的任务数为0并且上次记录的任务总数与当前的任务总数一致,考虑关闭集群if (taskRunning.get() == 0 && cacheTaskNum == taskNum.get()) {terminateCluster(cacheClusterId);MonitorLog.info("集群:" + cacheClusterId + "关闭成功,执行任务数为:" + taskNum.get());//关闭调度clusterWatchFuture.cancel(true);} else {//重新设置任务数cacheTaskNum = taskNum.get();}}, 11, 11, TimeUnit.MINUTES);}}/*** 初始化client操作*/protected abstract void init();/*** 关闭集群操作** @param clusterId clusterId*/protected abstract void terminateCluster(String clusterId);/*** 关闭client*/protected abstract void shutdown();/*** 发送创建集群的请求** @return 返回clusterId*/protected abstract String sendCreateRequest();/*** 判断以clusterPrefix开头的机器是否有存活** @return StringUtils.isBlank() 表示无存活*/protected abstract String getAliveId();/*** 检测集群是否创建完成** @param clusterId clusterId* @return*/protected abstract boolean isCompletion(String clusterId);/*** 获得master的ip** @param clusterId clusterId* @return*/protected abstract String getMasterIp(String clusterId);/*** 销毁集群,以及做一些后置操作*/protected synchronized void destroyCluster() {if (!clusterTerminate) {clusterTerminate = true;cacheIp = null;cacheClusterId = null;clusterWatchFuture = null;pool.shutdown();pool = null;shutdown();}}
}

阿里云集群的创建

这篇关于使用java在aliyun/aws创建E-MapReduce (emr)集群的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/694612

相关文章

SpringBoot多环境配置数据读取方式

《SpringBoot多环境配置数据读取方式》SpringBoot通过环境隔离机制,支持properties/yaml/yml多格式配置,结合@Value、Environment和@Configura... 目录一、多环境配置的核心思路二、3种配置文件格式详解2.1 properties格式(传统格式)1.

Apache Ignite 与 Spring Boot 集成详细指南

《ApacheIgnite与SpringBoot集成详细指南》ApacheIgnite官方指南详解如何通过SpringBootStarter扩展实现自动配置,支持厚/轻客户端模式,简化Ign... 目录 一、背景:为什么需要这个集成? 二、两种集成方式(对应两种客户端模型) 三、方式一:自动配置 Thick

使用Python构建智能BAT文件生成器的完美解决方案

《使用Python构建智能BAT文件生成器的完美解决方案》这篇文章主要为大家详细介绍了如何使用wxPython构建一个智能的BAT文件生成器,它不仅能够为Python脚本生成启动脚本,还提供了完整的文... 目录引言运行效果图项目背景与需求分析核心需求技术选型核心功能实现1. 数据库设计2. 界面布局设计3

使用IDEA部署Docker应用指南分享

《使用IDEA部署Docker应用指南分享》本文介绍了使用IDEA部署Docker应用的四步流程:创建Dockerfile、配置IDEADocker连接、设置运行调试环境、构建运行镜像,并强调需准备本... 目录一、创建 dockerfile 配置文件二、配置 IDEA 的 Docker 连接三、配置 Do

Spring WebClient从入门到精通

《SpringWebClient从入门到精通》本文详解SpringWebClient非阻塞响应式特性及优势,涵盖核心API、实战应用与性能优化,对比RestTemplate,为微服务通信提供高效解决... 目录一、WebClient 概述1.1 为什么选择 WebClient?1.2 WebClient 与

Android Paging 分页加载库使用实践

《AndroidPaging分页加载库使用实践》AndroidPaging库是Jetpack组件的一部分,它提供了一套完整的解决方案来处理大型数据集的分页加载,本文将深入探讨Paging库... 目录前言一、Paging 库概述二、Paging 3 核心组件1. PagingSource2. Pager3.

Java.lang.InterruptedException被中止异常的原因及解决方案

《Java.lang.InterruptedException被中止异常的原因及解决方案》Java.lang.InterruptedException是线程被中断时抛出的异常,用于协作停止执行,常见于... 目录报错问题报错原因解决方法Java.lang.InterruptedException 是 Jav

深入浅出SpringBoot WebSocket构建实时应用全面指南

《深入浅出SpringBootWebSocket构建实时应用全面指南》WebSocket是一种在单个TCP连接上进行全双工通信的协议,这篇文章主要为大家详细介绍了SpringBoot如何集成WebS... 目录前言为什么需要 WebSocketWebSocket 是什么Spring Boot 如何简化 We

java中pdf模版填充表单踩坑实战记录(itextPdf、openPdf、pdfbox)

《java中pdf模版填充表单踩坑实战记录(itextPdf、openPdf、pdfbox)》:本文主要介绍java中pdf模版填充表单踩坑的相关资料,OpenPDF、iText、PDFBox是三... 目录准备Pdf模版方法1:itextpdf7填充表单(1)加入依赖(2)代码(3)遇到的问题方法2:pd

Java Stream流之GroupBy的用法及应用场景

《JavaStream流之GroupBy的用法及应用场景》本教程将详细介绍如何在Java中使用Stream流的groupby方法,包括基本用法和一些常见的实际应用场景,感兴趣的朋友一起看看吧... 目录Java Stream流之GroupBy的用法1. 前言2. 基础概念什么是 GroupBy?Stream