Integrating scale out and fault tolerance in stream processing using operator state management

Integrating scale out and fault tolerance in stream processing using operator state management
复制标题

DOI:
10.1145/2463676.2465282
复制
发表时间:
2013-06
期刊:
--
影响因子:
--
通讯作者:
R. Fernandez;Matteo Migliavacca;Evangelia Kalyvianaki;P. Pietzuch
R. Fernandez;Matteo Migliavacca;Evangelia Kalyvianaki;P. Pietzuch
中科院分区:
其他
文献类型:
--
作者:
R. Fernandez;Matteo Migliavacca;Evangelia Kalyvianaki;P. Pietzuch

文献摘要

被引文献

相似文献

随着“大数据”应用程序的用户期望获得新的结果,我们见证了一种新的流处理系统(SPS),旨在扩展到大量的云托管机器。这些系统面临着新的挑战:(i)为了从云计算的“按需付费”模式中获益,它们必须按需扩展,在工作负载增加时获取额外的虚拟机(VM)并并行化操作员;(ii)在数百个VM上部署时,故障是常见的-系统必须具有快速恢复时间的容错性,但每台机器的开销较低。一个悬而未决的问题是,当流查询包括有状态操作符时,如何实现这两个目标,这些操作符必须在不影响查询结果的情况下进行扩展和恢复。我们的主要思想是通过一组状态管理原语将内部操作员状态显式地暴露给SPS。在此基础上,我们描述了一个集成的方法,动态扩展和恢复的状态运营商。外部化的运营商状态由SPS定期进行检查点设置,并备份到上游VM。SPS可识别各个运营商瓶颈,并通过分配新VM和划分检查点状态来自动扩展这些瓶颈。在任何时候,通过在新VM上恢复检查点状态并重播未处理的元组,可以恢复失败的操作符。我们使用Amazon EC2云平台上的Linear Road Benchmark评估了这种方法,并表明它可以自动扩展到负载因子L=350(50个虚拟机),同时从故障中快速恢复。
As users of "big data" applications expect fresh results, we witness a new breed of stream processing systems (SPS) that are designed to scale to large numbers of cloud-hosted machines. Such systems face new challenges: (i) to benefit from the "pay-as-you-go" model of cloud computing, they must scale out on demand, acquiring additional virtual machines (VMs) and parallelising operators when the workload increases; (ii) failures are common with deployments on hundreds of VMs-systems must be fault-tolerant with fast recovery times, yet low per-machine overheads. An open question is how to achieve these two goals when stream queries include stateful operators, which must be scaled out and recovered without affecting query results. Our key idea is to expose internal operator state explicitly to the SPS through a set of state management primitives. Based on them, we describe an integrated approach for dynamic scale out and recovery of stateful operators. Externalised operator state is checkpointed periodically by the SPS and backed up to upstream VMs. The SPS identifies individual operator bottlenecks and automatically scales them out by allocating new VMs and partitioning the checkpointed state. At any point, failed operators are recovered by restoring checkpointed state on a new VM and replaying unprocessed tuples. We evaluate this approach with the Linear Road Benchmark on the Amazon EC2 cloud platform and show that it can scale automatically to a load factor of L=350 with 50 VMs, while recovering quickly from failures.