并行计算(Parallel Computing)
并行计算或称平行计算是相对于串行计算来说的。并行计算(Parallel Computing)是指同时使用多种计算资源解决计算问题的过程。为执行并行计算,计算资源应包括一台配有多处理机(并行处理)的计算机、一个与网络相连的计算机专有编号,或者两者结合使用。并行计算的主要目的是快速解决大型且复杂的计算问题。
并行计算可以划分成时间并行和空间并行。时间并行即流水线技术,空间并行使用多个处理器执行并发计算,当前研究的主要是空间的并行问题。以程序和算法设计人员的角度看,并行计算又可分为数据并行和任务并行。数据并行把大的任务化解成若干个相同的子任务,处理起来比任务并行简单。
空间上的并行导致两类并行机的产生,按照Michael Flynn(费林分类法)的说法分为单指令流多数据流(SIMD)和多指令流多数据流(MIMD),而常用的串行机也称为单指令流单数据流(SISD)。MIMD类的机器又可分为常见的五类:并行向量处理机(PVP)、对称多处理机(SMP)、大规模并行处理机(MPP)、工作站机群(COW)、分布式共享存储处理机(DSM)。
分布式计算(Distributed Computing)
分布式计算这个研究领域,主要研究分散系统(Distributed system)如何进行计算。分散系统是一组计算机,通过计算机网络相互链接与通信后形成的系统。把需要进行大量计算的工程数据分区成小块,由多台计算机分别计算,在上传运算结果后,将结果统一合并得出数据结论的科学。
目前常见的分布式计算项目通常使用世界各地上千万志愿者计算机的闲置计算能力,通过互联网进行数据传输。如分析计算蛋白质的内部结构和相关药物的Folding@home项目,该项目结构庞大,需要惊人的计算量,由一台电脑计算是不可能完成的。即使现在有了计算能力超强的超级电脑,但是一些科研机构的经费却又十分有限。
分布式计算比起其它算法具有以下几个优点:
- 稀有资源可以共享。
- 通过分布式计算可以在多台计算机上平衡计算负载。
- 可以把程序放在最适合运行它的计算机上。其中,共享稀有资源和平衡负载是计算机分布式计算的核心思想之一。
并行计算与分布式计算的区别
简单的理解,引用Answers.com上一个答案:
Parallel computing and distributed computing are ways of exploiting parallelism in computing to achieve higher performance. Multiple processing elements are used to solve a problem, either to have it done faster or to have a larger size problem been solved. To state simply, if the processing elements share the memory, it is called parallel computing, otherwise it is called distributed computing. Some have opinion that distributed computing is a special form of parallel computing.
并行计算与分布式计算都是运用并行来获得更高性能,化大任务为小任务。简单说来,如果处理单元共享内存,就称为并行计算,反之就是分布式计算。也有人认为分布式计算是并行计算的一种特例。
但是分布式的任务包互相之间有独立性,上一个任务包的结果未返回或者是结果处理错误,对下一个任务包的处理几乎没有什么影响。因此,分布式的实时性要求不高,而且允许存在计算错误(因为每个计算任务给好几个参与者计算,上传结果到服务器后要比较结果,然后对结果差异大的进行验证。
分布式要处理的问题一般是基于寻找模式的。所谓的寻找,就相当于穷举法!为了尝试到每一个可能存在的结果,一般从0~N( 某一数值)被一个一个的测试,直到我们找到所要求的结果。事实上,为了易于一次性探测到正确的结果,我们假设结果是以某个特殊形式开始的。在这种类型的搜索里,我们也许幸运的一开始就找到答案;也许不够走运以至于到最后才找到答案,这都很公平。
这么说,并行程序并行处理的任务包之间有很大的联系,而且并行计算的每一个任务块都是必要的,没有浪费的分割的,就是每个任务包都要处理,而且计算结果相互影响,就要求每个的计算结果要绝对正确,而且在时间上要尽量做到同步,而分布式的很多任务块可以根本就不处理,有大量的无用数据块,所以说分布式计算的速度尽管很快,但是真正的效率是低之再低 的,可能一直在寻找,但是永远都找不到,也可能一开始就找到了;而并行处理不同,它的任务包个数相对有限,在一个有限的时间应该是可能完成的。
分布式的编写一般用的是C++(也有用JAVA的,但是都不是主流),基本不用MPI接口。并行计算用MPI或者OpenMP。
集群计算(Cluster Computing)
计算机集群将一组松散集成的计算机软件或硬件连接起来高度紧密地协作完成计算工作。在某种意义上,他们可以被看作是一台计算机。集群系统中的单个计算机通常称为节点,通常通过局域网连接,但也有其它的可能连接方式。集群计算机通常用来改进单个计算机的计算速度和/或可靠性。一般情况下集群计算机比单个计算机,比如工作站或超级计算机性价比要高得多。
根据组成集群系统的计算机之间体系结构是否相同,集群可分为同构与异构两种。集群计算机按功能和结构可以分为,高可用性集群(High-availability (HA) clusters)、负载均衡集群(Loadbalancing clusters)、高性能计算集群(High-performance (HPC)clusters)、网格计算(Grid computing)。
高可用性集群,一般是指当集群中有某个节点失效的情况下,其上的任务会自动转移到其他正常的节点上。还指可以将集群中的某节点进行离线维护再上线,该过程并不影响整个集群的运行。
负载均衡集群,负载均衡集群运行时,一般通过一个或者多个前端负载均衡器,将工作负载分发到后端的一组服务器上,从而达到整个系统的高性能和高可用性。这样的计算机集群有时也被称为服务器群(Server Farm)。一般高可用性集群和负载均衡集群会使用类似的技术,或同时具有高可用性与负载均衡的特点。Linux虚拟服务器(LVS)项目在Linux操作系统上提供了最常用的负载均衡软件。
高性能计算集群,高性能计算集群采用将计算任务分配到集群的不同计算节点儿提高计算能力,因而主要应用在科学计算领域。比较流行的HPC采用Linux操作系统和其它一些免费软件来完成并行运算。这一集群配置通常被称为Beowulf集群。这类集群通常运行特定的程序以发挥HPC cluster的并行能力。这类程序一般应用特定的运行库, 比如专为科学计算设计的MPI库。HPC集群特别适合于在计算中各计算节点之间发生大量数据通讯的计算作业,比如一个节点的中间结果或影响到其它节点计算结果的情况。
网格计算(Grid Computing)
网格计算是分布式计算的一种,也是一种与集群计算非常相关的技术。如果我们说某项工作是分布式的,那么,参与这项工作的一定不只是一台计算机,而是一个计算机网络,显然这种蚂蚁搬山的方式将具有很强的数据处理能力。网格计算的实质就是组合与共享资源并确保系统安全。
网格计算通过利用大量异构计算机的未用资源(CPU周 期和磁盘存储),将其作为嵌入在分布式电信基础设施中的一个虚拟的计算机集群,为解决大规模的计算问题提供一个模型。网格计算的焦点放在支持跨管理域计算 的能力,这使它与传统的计算机集群或传统的分布式计算相区别。网格计算的目标是解决对于任何单一的超级计算机来说仍然大得难以解决的问题,并同时保持解决 多个较小的问题的灵活性。这样,网格计算就提供了一个多用户环境。
集群计算与网格计算的区别
- 简单地,网格与传统集群的主要差别是网格是连接一组相关并不信任的计算机,它的运作更像一个计算公共设施而不是一个独立的计算机。网格通常比集群支持更多不同类型的计算机集合。
- 网格本质上就是动态的,集群包含的处理器和资源的数量通常都是静态的。在网格上,资源则可以动态出现,资源可以根据需要添加到网格中或从网格中删除。
- 网格天生就是在本地网、城域网或广域网上进行分布的。网格可以分布在任何地方。而集群物理上都包含在一个位置的相同地方,通常只是局域网互连。集群互连技 术可以产生非常低的网络延时,如果集群距离很远,这可能会导致产生很多问题。物理临近和网络延时限制了集群地域分布的能力,而网格由于动态特性,可以提供 很好的高可扩展性。
- 集群仅仅通过增加服务器满足增长的需求。然而,集群的服务器数量、以及由此导致的集群性能是有限的:互连网络容量。也就是说如果一味地想通过扩大规模来提高集群计算机的性能,它的性价比会相应下降,这意味着我们不可能无限制地扩大集群的规模。 而网格虚拟出空前的超级计算机,不受规模的限制,成为下一代Internet的发展方向。
- 集群和网格计算是相互补充的。很多网格都在自己管理的资源中采用了集群。实际上,网格用户可能并不清楚他的工作负载是在一个远程的集群上执行的。尽管网格与集群之间存在很多区别,但是这些区别使它们构成了一个非常重要的关系,因为集群在网格中总有一席之地 特定的问题通常都需要一些紧耦合的处理器来解决。然而,随着网络功能和带宽的发展,以前采用集群计算很难解决的问题现在可以使用网格计算技术解决了。理解网格固有的可扩展性和集群提供的紧耦合互连机制所带来的性能优势之间的平衡是非常重要的。
云计算(Cloud Computing)
云计算是最新开始的新概念,它不只是计算等计算机概念,还有运营服务等概念了。它是分布式计算、并行计算和网格计算的发展,或者说是这些概念的商业实现。云计算不但包括分布式计算还包括分布式存储和分布式缓存。分布式存储又包括分布式文件存储和分布式数据存储。
云计算与并行、分布式、网格和集群计算的区别
云计算是从集群技术发展而来,区别在于集群虽然把多台机器联了起来,但其某项具体任务执行的时候还是会被转发到某台服务器上,而云可以简单的认为是任务可以被分割成多个进程在多台服务器上并行计算,然后得到结果,好处在于大数据量的操作性能非常好。云可以使用廉价的PC服务器 ,可以管理大数据量与大集群,关键技术在于能够对云内的基础设施进行动态按需分配与管理。云计算与并行计算、分布式计算的区别,以计算机用户来说,并行计算是由单个用户完成的,分布式计算是由多个用户合作完成的,云计算是没有用户参与,而是交给网络另一端的服务器完成的。
计算模型
计算模型可以分为三类:顺序模型、函数模型和并发模型。
顺序模型
顺序模型包括:
- Finite state machines(有限状态机)
- Post machines (Post–Turing machines and tag machines).(后图灵机和标签机)
- Pushdown automata
- Register machines(注册机)
- Random-access machines(随机存取机器)
- Turing machines(图灵机)
- Decision tree model(决策树模型)
功能模型
功能模型包括:
- Abstract rewriting systems(摘要重写系统)
- Combinatory logic(组合逻辑)
- General recursive functions(一般递归函数)
- Lambda calculus(lambda演算)
并发模型
并发模型包括:
- Actor model
- Cellular automaton
- Interaction nets
- Kahn process networks
- Logic gates and digital circuits
- Petri nets
- Synchronous Data Flow
编程语言
Lambda calculus
Alonzo Church的lambda 演算可以被视为最早的消息传递编程语言(参见 Hewitt、Bishop 和 Steiger 1973;Abelson 和 Sussman 1985)。例如,下面的 lambda 表达式在为leftSubTree和rightSubTree提供参数时实现树数据结构。当这样一棵树被赋予参数消息“getLeft”时,它返回leftSubTree并且同样地当被赋予消息“getRight”时它返回rightSubTree。
λ(leftSubTree,rightSubTree)
λ(message)
if (message == "getLeft")then leftSubTree
else if (message == "getRight")then rightSubTree 然而,lambda 演算的语义是使用变量替换来表达的,其中参数的值被替换到调用的 lambda 表达式的主体中。替代模型不适合并发,因为它不允许共享不断变化的资源的能力。受 lambda 演算的启发,编程语言Lisp的解释器使用了一种称为环境的数据结构,因此不必将参数值代入调用的 lambda 表达式的主体中。这允许共享更新共享数据结构的效果,但不提供并发性。
Simula
Simula 67率先使用消息传递进行计算,受到离散事件模拟应用程序的推动。在以前的模拟语言中,这些应用程序变得庞大且非模块化。在每个时间步,一个大型中央程序必须检查并更新每个模拟对象的状态,这些状态根据它在该步骤中与之交互的任何模拟对象的状态而改变。 Kristen Nygaard和Ole-Johan Dahl提出了一个想法(在 1967 年的 IFIP 研讨会上首次描述),即在每个对象上都有方法,这些方法将根据来自其他对象的消息更新自己的本地状态。此外,他们还为对象引入了类结构继承。他们的创新大大提高了程序的模块化。
然而,Simula 使用协程控制结构而不是真正的并发。
Smalltalk
Alan Kay在开发Smalltalk -71时受到Planner模式导向调用中消息传递的影响。Hewitt 对 Smalltalk-71 很感兴趣,但被通信的复杂性推迟了,通信包括对许多领域的调用,包括global、sender、receiver、reply-style、status、reply、operator selector等。
1972 年,Kay 访问了麻省理工学院,并讨论了他对 Smalltalk-72 的一些想法,该想法建立在Seymour Papert的Logo作品和用于教孩子编程的“小人”计算模型的基础上。然而,Smalltalk-72 的消息传递相当复杂。语言中的代码被解释器视为简单的标记流。正如Dan Ingalls后来描述的那样:
在动态上下文中查找(在程序中)遇到的第一个(令牌),以确定后续消息的接收者。名称查找从当前激活的类字典开始。如果在那里失败,它会移动到该激活的发送者,依此类推发送者链。当最终找到令牌的绑定时,它的值成为新消息的接收者,解释器激活该对象类的代码。
因此,Smalltalk-72 中的消息传递模型与不适合并发的特定机器模型和编程语言语法紧密相关。此外,尽管系统是自举的,但语言结构并未正式定义为响应Eval消息的对象(参见下面的讨论)。这导致一些人相信基于消息传递的并发计算的新数学模型应该比 Smalltalk-72 更简单。
Smalltalk 语言的后续版本在很大程度上遵循了在程序的消息传递结构中使用 Simula虚拟方法的路径。但是Smalltalk-72把整数、浮点数等原语变成了对象。Simula 的作者曾考虑过将这些基元变成对象,但主要出于效率原因而没有考虑。 Java最初使用了整数、浮点数等原始版本和对象版本的权宜之计。C#编程语言(以及从 Java 1.5 开始的 Java 的更高版本)采用了使用装箱和拆箱的不太优雅的解决方案,它的一个变体在早期的一些Lisp实现中使用过。
Smalltalk 系统继续变得非常有影响力,在位图显示、个人计算、类浏览器界面和许多其他方面进行了创新。有关详细信息,请参阅 Kay 的The Early History of Smalltalk。[1] 与此同时,麻省理工学院的 Actor 工作仍然专注于开发更高级别并发的科学和工程。(请参阅 Jean-Pierre Briot 的论文,了解后来开发的关于如何将某些类型的 Actor 并发合并到 Smalltalk 的更高版本中的想法。)
Petri nets
在 Actor 模型开发之前,Petri 网被广泛用于对非确定性计算建模。然而,他们被广泛认为有一个重要的局限性:他们模拟控制流而不是数据流。因此,它们不容易组合,从而限制了它们的模块化。休伊特指出了 Petri 网的另一个难点:同步动作。 即,Petri 网中计算的原子步骤是一个转换,其中令牌同时从转换的输入位置消失并出现在输出位置。在他看来,使用具有这种同时性的原始人的物理基础是值得怀疑的。尽管存在这些明显的困难,Petri 网仍然是一种流行的并发建模方法,并且仍然是活跃研究的主题。
Threads, locks, and buffers (channels)
在 Actor 模型之前,并发是在线程、锁和缓冲区(通道)的低级机器术语中定义的。当然,Actor 模型的实现通常会使用这些硬件功能。但是,没有理由不能在不暴露任何硬件线程和锁的情况下直接在硬件中实现该模型。此外,计算中可能涉及的 Actors、线程和锁的数量之间没有必然关系。Actor 模型的实现可以以任何符合 Actor 法则的方式自由使用线程和锁。
分布式调度系统架构(定时任务架构)
不管是网络、内存、还是存储的分布式,它们最终目的都是为了实现计算的分布式:数据在各个计算机节点上流动,同时各个计算机节点都能以某种方式访问共享数据,最终分布式计算后的输出结果被持久化存储和输出。
分布式作为分布式系统里最重要的一个能力和目标,也是大数据系统的关技术之一。经过多年的发展与演进,目前业界已经存在很多成熟的分布式计算相关的开源编程框架和平台供我们选择。
单机定时任务
单机定时任务是最常见的,也是比较传统的任务执行方式,比如linux内置的Crontab。其通过cron表达式中分、时、日、月、周五种时间维度,实现单机定时任务的执行。
# 每晚的21:30重启smb
30 21 * * * /etc/init.d/smb restart 另外,在java中也有内置的定时任务,比如java.util.Timer类和它的升级版ScheduledThreadPoolExecutor,另外在Spring体系中也提供了Spring Task这种通过注解快速实现支持cron表达式的单机定时任务框架。
@EnableScheduling
@Service
public class ScheduledConsumerDemo {
@Value("${consumer.request.echoUrl}")
private String echoUrl;
/**
* 间隔1秒请求provider的echo接口
*
* @throws InterruptedException
*/
@Scheduled(fixedDelayString = "${consumer.auto.test.interval:1000}")
public void callProviderPer1Sec() throws InterruptedException {
String response = restTemplate.getForObject(echoUrl, String.class);
}
} 显而易见的,单机定时任务在应对简单的业务场景是很方便的,但在分布式架构已然成为趋势的现在,单机定时任务并不能满足企业级生产以及工业化场景的诉求,主要体现在集群任务配置统一管理、单点故障及单点性能、节点间任务的通讯及协调、任务执行数据汇总等方面。为了满足企业级生产的诉求,各类任务调度平台逐步兴起。
中心化调度(JAVA)
典型的中心化调度框架quartz,其作为任务调度界的前辈和带头大哥,通过优秀的调度能力、丰富的API接口、Spring集成性好等优点,使其一度成为任务调度的代名词。
quartz架构中使用数据库锁保障多节点任务执行时的唯一性,解决了单点故障的问题。但数据库锁的集中性也产生了严重的性能问题,比如大批量任务场景下,数据库成为了业务整体调度的性能瓶颈,同时在应用侧还会造成部分资源的等待闲置,另外还做不到任务的并行分片。
另一款出自大众点评的框架xxl-job,主要特点在于简单、易集成、有可视化控制台,相比quartz主要差异在于:
- 自研调度模块:
xxl-job将调度模块和任务模块解耦的异步化设计,解决了调度任务逻辑偏重时,调度系统性能大大降低的问题。其中,调度模块主要负责任务参数的解析及调用发起,任务模块则负责任务内容的执行,同时异步调度队列和异步执行队列的优化,使得有限的线程资源也可支撑一定量的job并发。
- 调度优化:
通过调度线程池、并行调度的方式,极大减小了调度阻塞的概率,同时提高了调度系统的承载量。
- 高可用保障:
调度中心的数据库中会保存任务信息、调度历史、调度日志、节点注册信息等,通过MySQL保证数据的持久化和高可用。任务节点的**故障转移(failover)**模式和心跳检测,也会动态感知每个执行节点的状态。
但由于xxl-job使用了跟quartz类似的数据库锁机制,所以同样不能避免数据库成为性能瓶颈以及中心化带来的其它问题。
去中心化调度(JAVA)
为了解决中心化调度存在的各种问题,国内开源框架也是八仙过海、尽显神通,比如口碑还不错的powerjob、当当的elastic-job、唯品会的saturn。saturn整体上是基于开源的elastic-job进行改进优化的,所以本文只针对powerjob和elastic-job做简要介绍。
powerjob诞生于2020年4月,其中包含了一些比较新的思路和元素,比如支持基于MapReduce的分布式计算、动态热加载Spring容器等。在功能上,多任务工作流编排、MapReduce执行模式、延迟执行是亮点,同时宣称所有组件都支持水平扩展,其核心组件说明如下:
- powerjob-server:调度中心,统一部署,负责任务调度和管理;
- powerjob-worker:执行器,提供单机执行、广播执行和分布式计算;
- powerjob-client:可选组件,OpenAPI客户端。
powerjob在解决中心化调度时的无锁调度设计思路值得借鉴,核心逻辑是通过appName作为业务应用分组的key,将powerjob-server和powerjob-worker以分组key进行逻辑绑定,即确保每个powerjob-worker集群在运行时只会连接到一台powerjob-server,这样就不需要锁机制来防止任务被多台server同时拿到,从而造成重复执行的问题。
虽然powerjob在各方面分析下来相对优秀,但毕竟产品迭代周期比较短,仍需要通过市场大规模应用来不断打磨产品细节,以验证产品的性能、易用性和稳定性。
elasticjob包含elasticjob-lite和elasticjob-cloud两个独立子项目,本文主要以elasticjob-lite为例展开。
elasticjob-lite定位为轻量级无中心化解决方案,在继承quartz的基础上,同时使用了zookeeper作为注册中心。在产品设计层面上,个人理解elasticjob相比其他分布式任务调度框架,更加侧重数据处理和计算,主要体现在如下两方面:
elasticjob-lite的无中心化:
- 没有调度中心的设计,在业务程序引入elasticjob的jar包后,由jar包进行任务的调度、状态通讯、日志落盘等操作。
- 每个任务节点间都是对等的,会在zookeeper中注册任务相关的信息(任务名称、对等实例列表、执行策略等),同时依赖zookeeper的选举机制进行执行实例的选举。
elasticjob-lite的弹性分片:
- 基于zookeeper,任务执行实例之间可以近乎实时的感知到对方的上下线状态,使得任务分片的分配可以随着任务实例数量的调整而调整,并且保证负载相对均匀。
- 在任务实例上下线时,并不会影响当前的任务,会在下次任务调度的时候重新分片,以避免任务的重复执行。
通过上述分析,elasticjob更多的是针对分布式任务计算场景设计,更适合做大量数据的分片计算或处理,尤其对资源利用率有要求的场景下更有优势。
其他框架(JAVA)
uncode-schedule
基于zookeeper,比较小众,不推荐
基于zookeeper+spring task/quartz的分布式任务调度组件,确保所有任务在集群中不重复,不遗漏的执行。支持动态添加和删除任务。
LTS
最近一次更新在2年前,目前不是很活跃;也比较小众,不推荐
LTS(light-task-scheduler)主要用于解决分布式任务调度问题,支持实时任务,定时任务和Cron任务。有较好的伸缩性,扩展性;
TBSchedule
阿里早期开源的分布式任务调度系统。代码略陈旧,使用timer而非线程池执行任务调度。众所周知,timer在处理异常状况时是有缺陷的。而且TBSchedule作业类型较为单一,只能是获取/处理数据一种模式。还有就是文档缺失比较严重。
依赖Zookeeper;
TBSchedule是一款非常优秀的高性能分布式调度框架,广泛应用于阿里巴巴、淘宝、支付宝、京东、聚美、汽车之家、国美等很多互联网企业的流程调度系统。tbschedule在时间调度方面虽然没有quartz强大,但是它支持分片功能。和quartz不同的是,tbschedule使用ZooKeeper来实现任务调度的高可用和分片。
TBSchedule的分布式机制是通过灵活的Sharding方式实现的,分片的规则由客户端决定,比如可以按所有数据的ID按10取模分片、按月份分片等等。TBSchedule的宿主服务器可以进行动态扩容和资源回收,这个特点主要是因为它后端依赖的ZooKeeper,这里的ZooKeeper对于TBSchedule来说是一个NoSQL,用于存储策略、任务、心跳信息数据,它的数据结构类似文件系统的目录结构,它的节点有临时节点、持久节点之分。调度引擎启动后,随着业务量数据量的增多,当前Cluster可能不能满足目前的处理需求,那么就需要增加服务器数量,一个新的服务器上线后会在ZooKeeper中创建一个代表当前服务器的一个唯一性路径(临时节点),并且新上线的服务器会和ZooKeeper保持长连接,当通信断开后,节点会自动摘除。
TBSchedule会定时扫描当前服务器的数量,重新进行任务分配。TBSchedule不仅提供了服务端的高性能调度服务,还提供了一个scheduleConsole的war包,随着宿主应用的部署直接部署到服务器,可以通过web的方式对调度的任务、策略进行监控管理,以及实时更新调整。
Saturn
Saturn是唯品会在github开源的一款分布式任务调度产品。它是基于当当elastic-job 1.0版本来开发的,其上完善了一些功能和添加了一些新的feature。
亮点:
支持多语言开发 python、Go、Shell、Java、Php。
管理控制台和数据统计分析更加完善
缺点:
技术文档较少 , 该框架是2016年由唯品会的研发团队基于elastic-job开发而来
Opencron
比较小众,不推荐
Antares
Antares 是一款基于 Quartz 机制的分布式任务调度管理平台,内部重写执行逻辑,一个任务仅会被服务器集群中的某个节点调度。用户可通过对任务预分片,有效提升任务执行效率;也可通过控制台 antares-tower 对任务进行基本操作,如触发,暂停,监控等。
Antares 是基于 Quartz 的分布式调度,支持分片,支持树形任务依赖,但是不是跨平台的
sia-task
没有仔细研究,依赖Zookeeper,感觉更偏向于解决编排和跨平台
无论是互联网应用或者企业级应用,都充斥着大量的批处理任务。我们常常需要一些任务调度系统帮助我们解决问题。随着微服务化架构的逐步演进,单体架构逐渐演变为分布式、微服务架构。在此的背景下,很多原先的任务调度平台已经不能满足业务系统的需求。于是出现了一些基于分布式的任务调度平台。这些平台各有其特点,但各有不足之处,比如不支持任务编排、与业务高耦合、不支持跨平台等问题。不是非常符合公司的需求,因此我们开发了微服务任务调度平台(SIA-TASK)。
SIA是我们公司基础开发平台Simple is Awesome的简称,SIA-TASK(微服务任务调度平台)是其中的一项重要产品,SIA-TASK契合当前微服务架构模式,具有跨平台,可编排,高可用,无侵入,一致性,异步并行,动态扩展,实时监控等特点。
并行计算模型
并行编程模型与计算模型密切相关。并行计算模型是用于分析计算过程成本的抽象,但它不一定需要实用,因为它可以在硬件和/或软件中有效地实现。相比之下,编程模型确实特别暗示了硬件和软件实现的实际考虑。
并行编程语言可以基于一种编程模型或其组合。例如,高性能Fortran基于共享内存交互和数据并行问题分解,Go提供了共享内存和消息传递交互的机制。
并行编程模型示例
| Name | Class of interaction | Class of decomposition | Example implementations |
| Actor model | Asynchronous message passing | Task | D, Erlang, Scala, SALSA |
| Bulk synchronous parallel | Shared memory | Task | Apache Giraph, Apache Hama, BSPlib |
| Communicating sequential processes | Synchronous message passing | Task | Ada, Occam, VerilogCSP, Go |
| Circuits | Message passing | Task | Verilog, VHDL |
| Dataflow | Message passing | Task | Lustre, TensorFlow, Apache Flink |
| Functional | Message passing | Task | Concurrent Haskell, Concurrent ML |
| LogP machine | Synchronous message passing | Not specified | None |
| Parallel random access machine | Shared memory | Data | Cilk, CUDA, OpenMP, Threading Building Blocks, XMTC |
| SPMD PGAS | Partitioned global address space | Data | Fortran 2008, Unified Parallel C, UPC++, SHMEM |
| Global-view Task parallelism | Partitioned global address space | Task | Chapel, X10 |
Actor
既被用作计算的理论理解框架,又被用作并发系统的几种实际实现的理论基础。
定义 Actor 模型的一个重要挑战是抽象出实现细节。
例如,考虑以下问题:“每个 Actor 是否都有一个队列,在该队列中存储其通信,直到被 Actor 接收并进行处理?” Carl Hewitt反对将此类队列作为 Actor 模型的组成部分。一个考虑因素是,此类队列本身可以建模为接收消息以对通信进行入队和出队的 Actor。另一个考虑因素是某些 Actor 在实际实现中不会使用此类队列。例如, Actor 可能有一个仲裁者网络。当然,有一个数学抽象是序列Actor 收到的通信。但这个序列仅在 Actor 操作时出现。事实上,此序列的排序可能是不确定的(请参阅并发计算中的不确定性)。
另一个抽象实现细节的例子是解释问题:“解释应该成为 Actor 模型的一个组成部分吗?” 解释的想法是,Actor 将由其程序脚本如何处理eval消息来定义。(通过这种方式,Actor 将以类似于Lisp的方式定义,Lisp 是由用 Lisp 编写的名为eval的元循环解释器过程“定义”的。)休伊特反对将解释集成到 Actor 模型中。一个考虑是处理eval消息,Actor 的程序脚本本身就会有一个程序脚本(反过来会有......)!另一个考虑是一些演员不会在他们的实际解释中使用解释。例如, Actor 可能会在硬件中实现。当然,解释本身并没有错。此外,使用eval消息实现解释器比 Lisp 的单一解释器方法更加模块化和可扩展。
终于在第一个Actor发表八年后,Will Clinger(在Irene Greif 1975、Gordon Plotkin 1976、Michael Smyth 1978、Henry Baker 1978、Francez、Hoare、Lehmann和de Roever 1979以及Milne和Milnor 1979的工作基础上)在1981年的论文中发表了第一个令人满意的数学指称模型,该模型使用领域理论纳入了无界的非确定性(见Clinger的模型)。随后,Hewitt[2006]用到达时间对图进行了扩充,构建了一个技术上更简单、更容易理解的指称模型。见指称语义学的历史。
Actor是计算机科学领域中的一个并行计算模型,它把Actor当做通用的并行计算原语:一个Actor对接收到的消息做出响应,进行本地决策,可以创建更多的Actor(子Actor),或者发送更多的消息;同时准备接收下一条消息。
在Actor理论中,一切都被认为是Actor,这和面向对象语言里一切都被看成对象很类似。但包括面向对象语言在内的软件通常是顺序执行的,而Actor模型本质上则是并发的。Actor之间仅通过发送消息进行通信,所有的操作都是异步的,不同的Actor可以同时处理各自的信息,使整个系统获得大规模的并发能力。
Actor模型简单原理图:
根据上图,每个Actor都有一个Mailbox(邮箱),Actor A 发送给消息给Actor B,就好像Actor A 给Actor B写了一封邮箱地址为Actor B的邮箱地址的邮件(消息)一样,随后平台负责投递邮件。当邮件Actor B之后,平台就会通知Actor B收取邮件并做出回复,如果有多封邮件,则Actor B按顺序处理。很简单和容易理解的技术,但是蕴含了强大的力量。
Actor B收到消息后可能会做那些处理呢?
1创建其他Actor、2向其他Actor发送消息、3指定下一条消息到来的行为,比如修改自己的状态
在什么情况下一个Actor会创建子Actor呢?
通常情况是为了并行计算,比如我们有10G的文件要分析处理,我们可以在根Actor里创建10个子Actor,让每个Actor分别处理一个文件,为此根Actor给每个子Actor发送一个消息,消息里包含分配给它的的文件编号(或位置),当子Actor完成处理后,就把处理好的结果封装为应答消息返回给根Actor,然后根Actor在进行最后的汇总与输出,下面是这个过程的示意图。
一个Actor与其所创建的Actor形成父子关系。在实际编程中,父Actor应该监督其所创建的子Actor的状态,原因是父Actor知道可能会出现那些失败情况,知道如何处理他们,比如重新产生一个新的子Actor 来重做失败的任务,或者某个Actor失败后就通知其他Actor终止任务。
通过上面对Actor模型原理的简单分析,我们来总结一下Actor模型的优缺点。
优点:
1)将消息收发、线程调度、处理竞争和同步的所有复杂逻辑都委托给了Actor框架本身,而且对应用来说是透明的,我们可以认为Actor只是一个实现了Runnable接口的对象。关注多线程并发问题时,只需要关注多个Actor之间的消息流即可。
2)符合Actor模型的程序很容易进行测试,因为任意一个Actor都可以被单独进行单元测试。如果测试案例覆盖了该Actor所能响应的所有类型的消息,我们就可以确定该Actor的代码十分可靠。
缺点:
1) Actor完全避免共享并且仅通过消息来进行交流,使得程序失去了精细化并发调控能力,所以不适合实施细粒度的并行且可能导致系统响应时延的增加。如果在Actor程序中引入一些并行框架,就可能会导致系统的不确定性。
2)尽管使用Actor模型的程序 比使用线程和锁模型的程序更容易调试,Actor模型仍会碰到死锁这一类的共性问题,也会碰到一些Actor模型独有的问题(例如信箱移溢出)。
Akka
Akka是一个源可用的工具包和运行时,可简化JVM上并发和分布式应用程序的构建。Akka 支持多种并发编程模型,但它强调基于 actor 的并发,其灵感来自于Erlang。
Java和Scala都存在语言绑定。Akka 是用 Scala 编写的,从 Scala 2.10 开始,Scala 标准库中的 actor 已被弃用,转而使用 Akka。
Akka虽然是Scala写成的,但是由于Scala最终还是编译为Java字节码运行在JVM上,所以我们可以认为Akka属于Java领域。
Akka处理并发的方法基于Actor模型。在Akka里,Actor之间通信的唯一机制就是消息传递。
Akka官方宣传是这样介绍Akka的:
- 对并发、并行程序的简单的高级别的抽象
- 异步、非阻塞、高性能的事件驱动编程模型
- 非常轻量级的事件驱动处理(1GB内存可容纳约270万个actors)
对于多线程编程,经常会出现的缺点:线程维护困难、子线程出错后难以恢复、线程阻塞时浪费时间和资源;另外对于密集的计算任务,我们的系统需要达到很高的并发性能,单机系统资源无法满足计算需求,而类似java的fork/join框架又很难在分布式环境下简易地对并发计算进行架构设计。
而akka框架非常适合解决上述问题。akka框架是一款高性能、高容错性的分布式并发应用框架。akka底层采用scala语言实现,并基于actor并发模型,天然拥有异步、分布式能力,且具有很好的并发性能和容错机制。
面对大量的计算任务,系统怎样才能快速实时地得到想要的结果呢?很显然,依靠单核CPU的处理能力已不足以进行如此密集的计算(摩尔定律的失效),一般情况下,我们的解决方案是:把计算拆分成多个子任务实现并行(单机多核或分布式集群)执行。
Akka从底层就解决了我们大多数分布式&并行程序常见的难题,让工程师更专注于业务实现,同时,它也保留了多个扩展接口及配置,便于满足个性化定制的需要
在akka中,整个actor体系被抽象成一个公共的actor系统,即ActorSystem,ActorSystem是一个层级结构,通过"父监督"模式和DeathWatch模式限定了actor的管理策略。
akka还提供了诸多的配套组件,例如网络服务、持久化等。
ActorSystem
actor组件
akka中actor组件具有几个特征:
- actor引用(Actor Reference) akka不能通过new的方式创建引用。代替new的方式是,通过actorOf(创建actor)或actorSelection(查找actor)等方式返回Actor对象引用。对开发者而言,actor位置透明(可能存在于本地或远程)。
- 状态(State) actor在不同时刻的状态通常用变量来标识。akka在底层为每个actor抽象一个轻量级的执行“线程”,实现了对状态的隔离。
- 行为(Behavior) actor在接受到消息后可以对消息进行处理,或者转发给其它actor处理。
- 监督策略 (Supervisor Strategy) actorSystem是一个层级结构,父actor对子actor具有监管的能力,可以针对子actor的异常进行:恢复、重启、停止、失败上溯等处理方案,另外提供了One-For-One(默认,只对异常的子actor进行处理)和All-For-One(对所有子actor进行处理)两种监督策略。
邮箱
每个actor都拥有自己的邮箱,所有接受的消息会先进入邮箱,actor从邮箱中取出消息进行处理。akka自带多种邮箱类型,并提供RequiresMessageQueue接口供开发者自定义特定类型的邮箱。
路由
actor消息也可以通过路由进行发送,路由可以是一个Router对象,也可以是一个自包含的actor:管理者所有的Routee(路由目标)。开发者可以根据需要选择不同的路由类型,如:轮询、随机、广播等。
网络服务
提供远程actor和分布式集群的基础能力,包含I/O、网络通信、序列化配置、gossip通信协议、节点管理、集群分片等
http、websocket模块
akka提供了处理 http、websocket协议的基础模块,可以在此基础上进行开发。
持久化
akka提供actor的状态持久化方案,在程序出错、宕机等场景下进行恢复。
目前Akka已经在多家互联网&软件公司广泛使用,比如eBay、Amazon、VMWare、PayPal、阿里、惠普、豌豆荚等,所涉行业包括游戏、金融投资、医疗保健、数据分析等。
使用场景包括:
- 服务后端,比如rest web,websocket服务,分布式消息处理等。
- 并发&并行,比如日志异步处理,密集数据计算等。
总之,对高并发和密集计算的系统,Akka都是适用的!
MapReduce
MapReduce[1]是Google提出的一个软件架构,用于大规模数据集的并行运算。概念“Map(映射)”和“Reduce(归约)”,及他们的主要思想,都是从函数式编程语言借鉴的,还有从矢量编程语言借来的特性。
当前的软件实现是指定一个“Map”(映射)函数,用来把一组键值对映射成一组新的键值对,指定并发的“Reduce”(归约)函数,用来保证所有映射的键值对中的每一个共享相同的键组。
Google在2004年发布三篇重要论文,也就是大数据的三驾马车:分布式文件系统GFS、计算引擎MapReduce、分布式数据库BigTable
Spark
Apache Spark 是一个用于大规模数据处理的开源统一分析引擎。Spark 为具有隐式数据并行性和容错性的集群编程提供了一个接口。
Spark 及其 RDD 于 2012 年开发,以应对MapReduce集群计算范式的局限性,该范式在分布式程序上强制采用特定的线性数据流结构:MapReduce 程序从磁盘读取输入数据,跨数据映射函数,减少映射,并将缩减结果存储在磁盘上。Spark 的 RDD 充当分布式程序的工作集,提供(故意)受限形式的分布式共享内存。[8]
在 Apache Spark 内部,工作流作为有向无环图(DAG) 进行管理。节点代表 RDD,而边代表 RDD 上的操作。
Spark 有助于实现迭代算法(在循环中多次访问其数据集)和交互式/探索性数据分析,即重复的数据库式数据查询。与Apache Hadoop MapReduce 实施相比,此类应用程序的延迟可能会减少几个数量级。[2] [9] 在迭代算法的类别中,有机器学习系统的训练算法,它形成了开发 Apache Spark 的最初动力。[10]
Apache Spark 需要一个集群管理器和一个分布式存储系统。对于集群管理,Spark支持standalone(原生Spark集群,可以手动启动集群,也可以使用安装包提供的启动脚本,也可以在单机上运行这些daemon进行测试)、Hadoop YARN、Apache Mesos或Kubernetes。[11]对于分布式存储,Spark 可以与多种接口,包括Alluxio、Hadoop 分布式文件系统(HDFS)、[12] MapR 文件系统(MapR-FS)、[13] Cassandra、[14] 可以实施OpenStack Swift、Amazon S3、Kudu、Lustre 文件系统[ 15]或自定义解决方案。Spark还支持一种伪分布式本地模式,通常只用于开发或测试目的,不需要分布式存储,可以使用本地文件系统代替;在这种情况下,Spark 运行在一台机器上,每个CPU 核心有一个执行程序。
Hadoop
当前热门的计算引擎都能归类到Hadoop的生态圈中,是08年以来一直都非常热门的技术。
最开始的Nutch是没有大规模数据存储和计算对应的解决方案,直到04年前后,Google发布了《The Google File System》、《MapReduce: Simplified Data Processing on Large Clusters 》和《Bigtable: A Distributed Storage System for Structured Data》,Nutch的作者看了之后用Java实现了Hadoop的三个核心组件:HDFS、MapReduce和HBase分别对应三篇论文阐述的理论思想。
Hadoop最开始是作为Nutch的一个子项目存在的,解决的问题就是搜索引擎对应的大规模数据存储、计算,直到被Yahoo孵化为一个完整的分布式系统解决方案并在08年之后全球广泛应用。
Hadoop作为一个分布式系统,也是由两部分组成的:分布式存储和分布式计算。MapReduce对应的就是最初Hadoop分布式计算的解决方案。
HDFS是分布式存储,MapRecue是分布式计算,对应的就是这两个模块,Hadoop2.x开始把资源管理单独拆分出来了,拆分出来的好处就是,YARN变成了一个公共的资源管理平台,在它上面不仅仅可以跑MapReduce程序,还可以跑很多其他的程序,只要你的程序满足YARN的规则即可
- HDFS负责海量数据的分布式存储
- MapReduce是一个计算模型,负责海量数据的分布式计算
- YARN主要负责集群资源的管理和调度
Hadoop发展到现在,可以集成在其中的计算框架非常多,从应用的角度讲大体上可以分为两种:离线和实时。
计算引擎按照出生年代、计算模型归了个类:
第一代:无实时计算能力
MapReduce
第二代:基于MR之上的工具与实时计算能力
Tez:基于MapReduce的DAG执行工具。
Storm:以拓扑结构的方式提供实时计算能力。
第三代:完整的技术栈与更加稳定的系统
Spark:提供更高效、更完整的技术栈。
Heron:比Storm更加稳定的实时系统。
第四代:更为先进的计算思想与框架的统一
Flink:最接近DataFlow数据流思想的解决方案。
Beam:多种计算框架的兼容与支持。
Storm
与Actor面向消息的分布式计算式模型不同,Apache Storm提供的是面向连续的消息流(Stream)的一种通用的分布式计算解决框架。
Apache Storm是一种侧重于极低延迟的流处理框架,也是要求近实时处理的工作负载的最佳选择。该技术可处理非常大量的数据,通过比其他解决方案更低的延迟提供结果。
实时数据计算系统Storm,它将计算数据的间隔缩短到足够小,收到一部分数据就计算一部分,也就是流式计算了。Storm填补了大数据实时计算的缺失,并且便于使用,一经推出就风靡业界。
Storm作为实时流式计算中的佼佼者,因其良好的特性使其使用场景非常广泛。
Zookeeper作为分布式协调服务框架,因其完善的数据一致性保证特性使其成为各框架必备组件。
1)日志处理: 监控系统中的事件日志,使用 Storm 检查每条日志信息,把符合匹配规则的消息保存到数据库。
2)电商商品推荐: 后台需要维护每个用户的兴趣点,主要基于用户的历史行为、查询、点击、地理信息等信息获得,其中有很多实时数据,可以使用 Storm 进行处理,在此基础上进行精准的商品推荐和放置广告。
Hadoop 是强大的大数据处理系统,但是在实时计算方面不够擅长;Storm的核心功能就是提供强大的实时处理能力,但没有涉及存储;所以 Storm 与 Hadoop 即不同也互补。
Storm与Hadoop应用场景对比:
Storm: 分布式实时计算,强调实时性,常用于实时性要求较高的地方
Hadoop:分布式批处理计算,强调批处理,常用于对已经在的大量数据挖掘、分析
Hadoop: 离线分析框架,适合离线的复杂的大数据处理
Spark:内存计算框架,适合在线、离线快速的大数据处理
Storm: 流式计算框架,适合在线的实时的大数据处理
Ray:系统中有三个基础概念Task、Object、Actor,分别对应分布式任务、对象和服务。而常用的面向过程的编程语言中,也刚好有三个基本概念,函数、变量和类。在Ray系统中,可以通过简单的改动,实现它们之间的转换。
单机程序 + Ray => 分布式计算程序
流计算 + 图计算 + Ray => 流图融合
流计算 + 机器学习 + Ray => 在线学习
分布式计算原理、发展、瓶颈
分布式计算(Distributed Computing)也可译为分散式运算,它主要研究如何应用分布式系统(Distributed System)进行计算。分布式系统中的组件位于不同的计算机上,它们之间通过消息传递进行交流、协作,最终实现一个共同的目标。组件之间的并发行、没有全局时钟、组件的独立故障是分布式系统中的三个主要特性。从基于SOA的系统到大型多人在线游戏,再到P2P都是分布式系统的应用。
在分布式系统中运行的计算机程序被称为分布式程序(distributed program)。在分布式系统中,实现消息传递的机制有很多,比如HTTP、类RPC连接器、MOM(Message-oriented middleware)等。
分布式计算也可应用于解决计算问题中。在分布式计算中,一个问题被分解为很多不同的子问题/任务,每一个任务再由一台或多台计算机解决。
要实现分布式计算首先要解决其中两个最重要的问题:
- 1.如何拆分计算逻辑
- 2.如何分发计算逻辑
拆分逻辑
计算逻辑要实现分布式,就必须要解决:如何将一个巨大的问题拆分成相对独立的子问题分发到各个机器上求解。
从在哪里发生计算的角度来看,所有的计算逻辑都能够划分为这两种类型:
- 1.能够分发到各个节点上并行执行的
- 2.需要经过一定量的结果合并之后才能继续执行的
两者之间需要建立一种可靠的通讯机制以保证准确无误地完成计算任务。
为了使整个求解过程完美的衔接起来,你需要解决一系列的通信、容灾、任务调度等问题。
首先对此公开提出解决方案的是Google的MapReduce。
其将一个分布式任务定义为两种类型的Job组成:
- 1.Map Job对应的就是可以在各个节点上一起执行互不影响的逻辑
- 2.Reduce Job处理的则是Map产生的中间结果(如果有)
Map和Reduce之间通过一个Shuffle过程来链接。
组成一个完整的分布式计算流程所有的细节问题都将在这个过程中得到有效的处理。
分发逻辑
计算逻辑拆分的问题解决之后,面临的将是如何将拆开的逻辑分发出去。这里将是分布式计算和集中式计算最大的不同点,移动计算逻辑而不移动数据。
对于Map类型的任务其实可以归为分布式存储系统的问题。因为基本上所有分布式计算框架都是基于优先数据本地化进行的。也就是说数据存在哪,计算就分发到哪。
分布式存储系统中海量数据的分发方案:
- 1.数据量分布
- 2.数据范围分布
- 3.Hash分片/消息队列
- 4.一致性Hash
对于Reduce类型的任务分发的重点取决于之前的Map和Shuffle过程。
分布式编程通常采用以下几种基本的框架:
- 客户服务器模式(Client-server model):把客户端 (Client) 与服务器 (Server) 区分开来。每一个客户端软件的实例都可以向一个服务器或应用程序服务器发出请求;
- 三层架构(Three-tier):将客户端移到中间层,无状态客户端可被使用。这种架构使应用时的部署变得简单,大部分的网页应用程序都是基于这种架构的;
- 多层架构(n-tier):多层架构是开发人员在开发过程当中面对复杂且易变的需求采取的一种以隔离控制为主的应对策略。每一层都可以单独部署。将整个项目自下而上的分为:数据持久(数据访问)层,逻辑(业务)层,UI(展现)层;
- 对等网络(P2P):是无中心服务器、依靠用户群(peers)交换信息的互联网体系,它的作用在于,减低以往网路传输中的节点,以降低资料遗失的风险。
典型分布式计算技术
- 中间件(Middleware)技术:属于可复用软件的范畴,处于操作系统软件与用户应用软件中间。中间件在操作系统、网络和数据库之上、应用软件之下,其作用是为处于上层的应用软件提供运行与开发的环境,帮助用户灵活、高效地开发和集成复杂的应用软件。
- 移动Agent技术:移动Agent是一个能在异构网络中自主地从一台主机迁移到另一台主机、并可与其他agent或资源交互的程序。移动Agent具有自治性、移动性、智能性。
- 网格(Grid):网格技术区别于传统的集中式大规模资源共享、分布式计算以及高性能计算等技术。它在动态的一组个体、机构和资源的虚拟组织中实行灵活、可靠、可调整的资源共享环境。在此环境中,网格所需解决的问题包括:唯一性认证、资源访问、资源发现的方式等。网格供应商向用户提供高性能计算环境,而信息系统无需购买昂贵的计算设备,只需从网格中获取所需的计算能力。
- Web Service技术:是对象/组件技术在 Interne中t 的延伸,是一种部署在 Web 上 的对象/组件。Web Service结合了以组件为基础的开发模式以 及 Web 的出色性能,一方面,Web Service和组件一样,具有黑 匣子的功能,可以在不关心功能如何实现的情况下重用;同时, 与传统的组件技术不同,Web Service可以把不同平台开发的 不同类型的功能块集成在一起,提供相互之间的互操作。
- P2P技术:P2P 系统由若干互联协作的计算机构成,是Internet上实 施分布式计算的新模式。它把C/S与B/S系统中的角色一体化, 引导网络计算模式从集中式向分布式偏移,也就是说网络应 用的核心从中央服务器向网络边缘的终端设备扩散,通过服 务器与服务器、服务器与PC机、PC机与PC机、PC机与WAP 手机等两者之间的直接交换而达成计算机资源与信息共享。
Web Service 技术的体系结构与基于中间件分布式系统 的体系结构相比,发现它们是非常相似的,可以把体系结构中 的 Web 程序看作中间件。从结构上来看,Web 服务只是从侧 面对中间件平台技术进行革新,虽然所有服务之间的通信都 以 XML 格式的消息为基础,但调用服务的基本途径主要还是RPC,而且具体实现并没有提供一种全新的编程模式。
网格技术与基于中间件的分布式计算技术相比较,它依 然以“中间件”为技术核心,在实现形式上并没有太大的改变。 然而经过一系列的技术革新,网格系统中的技术内涵已经发 生了深刻的变化。其一,基于中间件的分布式计算技术的资源 主要是指数据和软件,而网格计算的资源已经延伸到所有用 于共享的实体,包括硬件、软件,甚至分布式文件系统、缓冲池等;其二,在Internet上,网格中间件层提供了与Web服务一样优 秀的扩展功能,打破了传统分布式技术C/S模式的局限。
网格计算、Web Service等技术在异构平台上构筑了一层 通用的、与平台无关的信息和服务交换设施,从而屏蔽了 Internet中千差万别的差异,使信息和服务畅通无阻地在计算 机之间流动。网格计算与Web Service技术的共同载体是 Internet。但两者的不同之处在于,网格系统连接物理上分散的 硬件资源,形成虚拟计算组织,从而使计算资源得到充分共享。 而Web服务则是以商务应用为背景,是基于网格系统之上的。 网格系统为Web服务提供一个与硬件无关的虚拟计算机;而 Web服务是架构在虚拟计算机平台上,与环境、语言无关的应 用集成平台。
尽管各种分布式计算技术在理念、规范和实现等方面有 较大的差异,但它们之间并不矛盾,而是一种承上启下的关系, 有时甚至是融合的。因此,各种分布式计算技术可以共同存在, 它们的相互结合也是非常有意义和现实的
主要事件
| 年份 | 事件 | 相关论文/Reference |
| 1990年左右 | 第一代分布式计算技术诞生,开放式软件基金OSF(Open Software Foundation)定义了一个独立于操作系统和体系结构,用于开发分布式应用程序的环境OSF/DEC。 | |
| 1995年左右 | 第二代分布式计算技术大约诞生,发展迅速并逐渐成熟。面向对象思想和分布式计算技术的有机结合形成了分布式对象计算模型。OMG针对面向过程分布式计算技术的缺点,提出了OMA(Object Management Architecture参考模型),采用了分布式对象技术,其核心技术为ORB(Object Request Broker)。ORB如同一条总线把分布式对象系统中的各类对象和应用连接成相互作用的整体。 | |
| 2004 | MapReduce是一种分布式计算模型,是Google提出的,主要用于搜索领域,解决海量数据的计算问题。 | "Our abstraction is inspired by the map and reduce primitives present in Lisp and many other functional languages." -"MapReduce: Simplified Data Processing on Large Clusters", by Jeffrey Dean and Sanjay Ghemawat; from Google Research |
| 2014 | Apache Spark 是专为大规模数据处理而设计的快速通用的计算引擎。Spark是UC Berkeley AMP lab (加州大学伯克利分校的AMP实验室)所开源的类Hadoop MapReduce的通用并行框架,Spark,拥有Hadoop MapReduce所具有的优点;但不同于MapReduce的是——Job中间输出结果可以保存在内存中,从而不再需要读写HDFS,因此Spark能更好地适用于数据挖掘与机器学习等需要迭代的MapReduce的算法。 | X. Meng, J. Bradley, B. Yuvaz, E. Sparks, S. Venkataraman, D. Liu, J. Freeman, D. Tsai, M. Amde, S. Owen, D. Xin, R. Xin, M. Franklin, R. Zadeh, M. Zaharia, and A. Talwalkar. MLlib: Machine Learning in Apache Spark, JMLR, 17(34):1–7, 2016. |
| 2015 | TensorFlow最初由谷歌大脑团队开发,用于Google的研究和生产,于2015年11月9日在Apache 2.0开源许可证下发布。TensorFlow是一个开源软件库,用于各种感知和语言理解任务的机器学习。 | Dean, Jeff; Monga, Rajat; et al. (November 9, 2015). "TensorFlow: Large-scale machine learning on heterogeneous systems" (PDF). TensorFlow.org. Google Research. Retrieved November 10, 2015. |
| 2016 | MXNet由dmlc(Distributed (Deep) Machine Learning Community)打造,其大部分成员是中国人,以陈天奇,李沐,解浚源等为代表,创建了这个世界上目前排名第四的深度学习框架:MXNet。 | Chen, T., Li, M., Li, Y., Lin, M., Wang, N., Wang, M., Xiao, T., Xu, B., Zhang, C., & Zhang, Z. (2015). MXNet: A Flexible and Efficient Machine Learning Library for Heterogeneous Distributed Systems. CoRR, abs/1512.01274. |
| 2016 | PyTorch是使用GPU和CPU优化的深度学习张量库。于2016年首次发布。 | "An Introduction to PyTorch – A Simple yet Powerful Deep Learning Library". analyticsvidhya.com. Retrieved 2018-06-11. |
瓶颈
分布式计算技术是计算机网络的产物,也是计算机网络应用未来的发展方向,虽然已经产生了大量的技术产品,但同时也出现了一些问题:
- 标准问题:目前几乎所有的分布式计算技术都没有完整的统一标准,标准的缺乏使得资源的共享和社会化程度低下。分布式计算之间缺乏有效的交互、协作与协同,从而在很大程度上制约了其发展。
- 技术问题:分布式计算技术是分布式计算能够转化为工业化生产能力的前提,也是分布式计算的质量、功能的保证。可是到目前为止,所有的分布式计算技术都或多或少地存在没有解决的问题,没有哪一种技术能实现完全意义上的分布式计算,满足所有分布式计算的需求。
- 异构问题:现在的网络是一个异构的环境,分布式计算技术首先需要解决异构环境下的互操作问题。首要任务是如何解决互相识别。目前,不可能实现对所有的资源采用同一种描述方式,也无法智能地识别这些资源,这就导致分布式技术技术只能在一定的范围内使用。
- 安全性问题:分布式计算系统存在如下的安全问题:身份认证、数据传输安全、主机资源保护问题、恶意主机问题。为了保证分布式计算的安全,传统的安全机制不是特别有效,表现在分布式计算中事物交互的任一方都无法判断对方是否真正的可信,如:请求是否是病毒和木马所发起;远程计算机的环境是否存在信息泄漏或者恶意欺骗等。
未来发展方向
要创建可信的、松散的、健壮的分布式系统,使分布式计算得到真实的应用和普及,必须加快标准的研制,从底层信号的传输到复杂业务的流程等各种不同的层次都要形成统一的标准。此外,应该尽可能对分布式计算技术进行有效的组合,集中各种技术的优势,满足分布式计算的需求,这应该是分布式计算技术研究的方向。
以下列举几个分布式计算重点发展的方向:
可靠性,分布式系统由大量分布式组件构成。这些组件的数量和异构性,大大增加了分布式计算系统出错的可能。如何使得系统在部分组件出错的情况下仍然维持正常工作,是分布式计算一直以来的热点议题。
低延时,如今大量的计算任务需要在极短的时间内完成,因为用户的体验往往取决于计算服务的响应速度。分布式计算机系统本身的架构意味着组件间通信需要很大的时间开销。如何规避这些时间开销,获得低延时的计算响应速度对于当前以及外来分布式系统的应用意义重大。
一致性,分布式计算的组件各自独立运算,如何保证组件间的数据具有一致性,同时不引入大量的通信开销,也是未来的研究方向之一。
分布式系统分类
有几个大的维度来区分:
- 有状态、无状态
- 重存储还是重计算
- long service还是批处理。
一些常见的分布式系统大类:
- 支持持久化存储的分布式存储系统
- 着重计算的分布式/并行计算框架
- 分布式消息队列
根据不同的应用的领域,把上述分类细化,常见分布式存储系统分为:
- 分布式协同系统(分布式日志复制)
- 分布式任务调度框架
- 流计算框架
- 分布式文件/对象系统
- 分布式NoSQL存储
- 分布式关系数据库(OLAP、OLTP);
- 各种消息队列mq
- 分布式机器学习/深度学习训练框架
- 分布式协调系统(日志复制系统)其实就是paxos算法及其变体的实现,典型的有zookeeper、etcd;一般来说只存少量的元数据信息,重点在高可用强一致,不提供高的through put,是分布式系统不可或缺的组件;
- 面向非结构化数据的分布式文件/对象系统比较有名的包括Lustre(HPC)GlusterFS(NAS NFS)、HDFS(hadoop)、ceph(虚机块存储)、swift(restful对象存储),各有不同的适用领域。
- 结构化数据的NoSQL分布式存储,种类和数量最多,按照Martin Fowler的分类,包括Aggregated Oriented NoSQL和图数据库NoSql;Aggregated Oriented NoSQL大致分为3类:
- Key-value NoSQL,例如Redis Riak等;
- column family NoSQL(wide column store),典型的是Hbase Cassandra
- document NoSQL,典型的是MongoDB
开源的图数据库有Neo4j;
- 很多OLAP分析型数据库是分布式的,例如GreenPlum, CitusDB(都是基于pgSQL);还有一些基于proxy实现关系数据库集群(很多互联网公司都有类似mysql fabric的产品),当然这些并非真正意义上的分布式数据库,可以认为是个sql的分库分表中间件;
- 还有一类兼容单机数据库,支持强事务的OLTP分布式数据库,统称NewSQL;NewSQL有不同的类别,比如stone braker搞的voltdb,有类似google的Spanner/F1的tidb/cockroachDB,具体对比可以看下CMU Pavlo的newSQL综述文章。
- 开源的消息队列也非常多,应用广泛程度不亚于nosql存储。有些对事物支持比较好的消息队列,例如rabbitmq active mq等;还有很多的kafka,主要做日志处理。消息队列的核心关注点就是一个消息at least once, at most once, only once。
- 分布式计算 广义的分布式计算的外延很大,可以参考Distributed Computing Paradigms(包含了rpc p2p mpi等分布式通信模型)。AI/Big data兴起后,分布式并行计算的框架开始流行,比如
- hadoop/spark等以MapReduce为范式的批处理引擎
- flink,storm等流计算/CEP引擎
- tensorflow mxnet等ParameterServer范式的DeepLearning 训练框架
- 分布式图计算引擎比如pregel等
- 面向大数据集的传统ML的分布式训练框架,xgboost/lightGBM/gensim/Angel
Docker
Docker 是一个软件平台,让您可以快速构建、测试和部署应用程序。Docker 将软件打包成名为容器的标准化单元,这些单元具有运行软件所需的所有功能,包括库、系统工具、代码和运行时。程序在这个容器里运行,就好像在真实的物理机上运行一样。有了 Docker,就不用担心环境问题。
总体来说,Docker 的接口相当简单,用户可以方便地创建和使用容器,把自己的应用放入容器。容器还可以进行版本管理、复制、分享、修改,就像管理普通的代码一样。
Docker 的主要用途,目前有三大类。
(1)提供一次性的环境。比如,本地测试他人的软件、持续集成的时候提供单元测试和构建的环境。
(2)提供弹性的云服务。因为 Docker 容器可以随开随关,很适合动态扩容和缩容。
(3)组建微服务架构。通过多个容器,一台机器可以跑多个服务,因此在本机就可以模拟出微服务架构。
Docker技术的三大核心概念,分别是:
- 镜像(Image)
- 容器(Container)
- 仓库(Repository)
kubernetes
K8S,就是基于容器的集群管理平台,它的全称,是kubernetes。
一个K8S系统,通常称为一个K8S集群(Cluster)。
这个集群主要包括两个部分:
- 一个Master节点(主节点)
- 一群Node节点(计算节点)
一看就明白:Master节点主要还是负责管理和控制。Node节点是工作负载节点,里面是具体的容器。
Master节点包括API Server、Scheduler、Controller manager、etcd。
API Server是整个系统的对外接口,供客户端和其它组件调用,相当于“营业厅”。
Scheduler负责对集群内部的资源进行调度,相当于“调度室”。
Controller manager负责管理控制器,相当于“大总管”。
Node节点包括Docker、kubelet、kube-proxy、Fluentd、kube-dns(可选),还有就是Pod。
Pod是Kubernetes最基本的操作单元。一个Pod代表着集群中运行的一个进程,它内部封装了一个或多个紧密相关的容器。除了Pod之外,K8S还有一个Service的概念,一个Service可以看作一组提供相同服务的Pod的对外访问接口。
前几年,大家以为虚拟机是核心网的终极形态。目前看来,更有可能是容器化。这几年经常说的NFV(网元功能虚拟化),也有可能改口为NFC(网元功能容器化)。
以VoLTE为例,如果按以前2G/3G的方式,那需要大量的专用设备,分别充当EPC和IMS的不同网元。
而采用容器之后,很可能只需要一台服务器,创建十几个容器,用不同的容器,来分别运行不同网元的服务程序,这些容器,随时可以创建,也可以随时销毁。还能够在不停机的情况下,随意变大,随意变小,随意变强,随意变弱,在性能和功耗之间动态平衡。
Volcano
Volcano 是用于在 Kubernetes 上运行高性能工作负载的系统。它具有 Kubernetes 无法提供但许多类别的高性能工作负载通常需要的强大批处理调度功能,包括:
- 机器学习/深度学习
- 生物信息学/基因组学
- 其他大数据应用
这些类型的应用程序通常在 TensorFlow、Spark、PyTorch 和 MPI 等通用域框架上运行。
Volcano 主要用于AI、大数据、基因、渲染等诸多高性能计算场景,对主流通用计算框架均有很好的支持。它提供高性能计算任务调度,异构设备管理,任务运行时管理等能力,目前在很多领域都已落地应用。
Volcano主要是基于Kubernetes做的一个批处理系统,希望上层的HPC、中间层大数据的应用以及最下面一层AI能够在统一Kubernetes上面运行的更高效。
Volcano要解决的问题
面向高性能负载的调度策略
支持多种作业生命周期管理
支持多种异构硬件
面向高性能负载的性能优化
支持资源管理及分时共享