S4分布式流计算引擎

简介:

背景

  最近花了点时间研究了下分布式计算这一块的内容。领导给的第一个任务,就是学习下S4和GridGain。花了几天的时间把s4的源码看了下,把自己的理解和学习的内容做一个记录。 下一篇会是GridGain的分享

学习

s4是什么?
1.  s4的全称 :  Simple Scalable Streaming System (简单的描述:分布式流计算系统)
2.  特点: 
  • distributed(分布式)
  • scalable(扩展性)
  • partially fault-toleran(部分容错性)
  • pluggable (可插拔)
3. 产生的原因: 
  • Yahoo发起,主要用于解决"cost-per-click“广告,通过实时计算预测用户对广告的可能的点击行为。
  • 不用hadoop的原因: hadoop主要解决batch处理,基于mapReduce对可控的数据的进行处理。而流计算是针对不可控的点击事件,对实时性有严格要求。
4.  适用的场景:
  • 业务允许部分容错性。 (s4没有严格的failover机制,运行节点突然crash时,会导致当前节点中的数据丢失。后续的请求会failover到其他的节点上)

 

S4的设计:
 

容器概念(http://docs.s4.io/manual/overview.html)

  1. PE : Processing Elements (处理节点)
    * emit one or more events which may be consumed by other PEs,
    * publish results, possibly to an external data store or consumer.
  2. Events :  message   (消息)
    * arbitrary Java Objects
    * passed between PEs. (send and receive)
  3. PEC : processing element container  (处理节点容器)
    * invokes the appropriate PEs in the appropriate order
  4. node : a Processing Endpoint  (机器节点)
    * a jvm instance
    * contains a PEC
  5. cluster: a group nodes  (节点集群)

说明:

  • 一个cluster包含多个node
  • 每个node包含一个PE Container
  • 每个PE Container 包含多个PE
  • 每个PE消费Events,生成新的Events并传递给下一个PE

 

总体结构图:

  1. PE Container/PE
  2. EventListenter
  3. Dispatcher, EventEmitter
  4. Communcation

 

 

 

 

PE内部概念:(4个部分组成)

  • its functionality as defined by a PE class and associated configuration,
  • the named stream that it consumes,
  • the keyed attribute in those events, and
  • the value of the keyed attribute in events which it consumers

PE类关系图:



说明:

  • Persister : 每个PE对应一个Persister,s4中每个PE对应于一个key的value结果。每个value结构都将作为PE的一个instance
  • FrequenceType :  每个PE会定期进行flush output输出,可选择的策略(定时,数量阀值)
  • Clock : 每个PE的时间控制单元,有几种时间。(WallClock:基于系统时间处理 , EventClock:基于event事件时间控制)

重点理解一下: Keyless PE概念和PE Prototype

  • PE在底层实现了会以多实例存在。存储的key即为其keyed对应的value值,内部有个lookupTable概念。
  • 针对Keyless PE,其对应的存储key即为"*",所以每次通过lookupTable.get(value)返回的即为同一个节点,单例化
  • 针对prototype,其对应的存储就为其value,所以每次会根据当前keyed attribute确定返回的PE节点,基于这点可以实现PE节点数据的Join处理

 

EventListener/EventEmitter:

 


 说明: 

  • 每个PE Container是一个EventProducer,使用CommLayerListener做为其事件处理器,处理EventWrapper反序列化。
  • 每个PE包含一个Dispatcher,Dispatcher里包含了一个EventEmitter处理对应的EventWrapper对象的发送
  • 底层实际通讯的类: SenderProcess/ListenerProcess

 

Dispatcher类关系图:



 

说明: 

  • Partitioner, 每个dispatcher针对发送的目标cluster,会根据对应的key进行分区处理,路由到其中的node节点。(node节点的信息可以通过zookeeper进行动态管理)

 

考虑集群node节点的管理(node的新增 or 修改)

说明:

  • ProcessMonitor 监控当前运行node节点的状态,有static/dymaic两种维护状态
  • TaskManager  创建node节点,主要是设置lock文件,有static/dymaic两种维护状态

运行(PE状态变化)

 

S4缺点:

S4产品还是一个半成品,整体代码结构组织和风格上还是比较乱的,选择使用时需谨慎。存在的一些问题:

  1. failover (运行node节点出现crash,当前node上的PE数据将无法实施failover)
  2. persist (目前支持方式过于简单,需要考虑网络持久化,类似于nfs,分布式文件系统等,配合failover机制)
  3. communication  (只支持udp协议,数据传输可靠性上)
  4. load balancer (根据系统负载进行智能LB,目前暂时未看到相关实现。系统运行分为两种模式static or dymaic模式, static不存在智能调节LB处理)
  5. deploy  (手工方式介入deploy,无法支持apps的zero deploy模式。系统分为cluster/node两概念,node对应于一计算节点实例,cluster为一组处理相同业务的计算节点)
相关文章
|
消息中间件 算法 Java
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的保障容量的三大关键方案实现
尽管经过了上一篇文章 《【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的低延迟可用性机制方案实现》有了低延迟的优化保障,消息引擎仍需精心规划其容量。为了提供无与伦比的流畅体验,消息引擎必须实施有效的容量管理策略。
364 2
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的保障容量的三大关键方案实现
|
存储 边缘计算 人工智能
云计算与分布式系统架构:驱动数字化时代的创新引擎
本文将探讨云计算与分布式系统架构在数字化时代中的重要性,介绍其基本概念和原理,并探讨其在推动技术创新、提升企业效率和满足用户需求方面的作用。同时,还将提出未来发展的趋势和挑战,为读者提供对云计算与分布式系统架构的深入理解。
|
SQL 分布式计算 数据库连接
大数据Spark分布式SQL引擎
大数据Spark分布式SQL引擎
566 0
|
消息中间件 存储 负载均衡
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的HA高可用解决方案
昔之善战者,先为不可胜,以待敌之可胜。不可胜在己,可胜在敌。故善战者,能为不可胜,不能使敌之必可胜。故曰:胜可知,而不可为。
604 2
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的HA高可用解决方案
|
11月前
|
消息中间件 分布式计算 资源调度
《聊聊分布式》ZooKeeper与ZAB协议:分布式协调的核心引擎
ZooKeeper是一个开源的分布式协调服务,基于ZAB协议实现数据一致性,提供分布式锁、配置管理、领导者选举等核心功能,具有高可用、强一致和简单易用的特点,广泛应用于Kafka、Hadoop等大型分布式系统中。
|
消息中间件 存储 Java
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的低延迟可用性机制方案实现
在充满挑战的2023年度,我们不可避免地面对了一系列棘手的问题,例如响应速度缓慢、系统陷入雪崩状态、用户遭受不佳的体验以及交易量的下滑。这些问题的出现,严重影响了我们的业务运行和用户满意度,为了应对这些问题,我们所在团队进行了大量的研究和实践,提出了低延迟高可用的解决方案,并在分布式存储领域广泛应用。
276 2
【亿级数据专题】「分布式消息引擎」 盘点本年度我们探索服务的低延迟可用性机制方案实现
|
存储 分布式计算 分布式数据库
【专栏】云计算与分布式系统架构在数字化时代的关键作用。云计算,凭借弹性、可扩展性和高可用性,提供便捷的计算环境
【4月更文挑战第27天】本文探讨了云计算与分布式系统架构在数字化时代的关键作用。云计算,凭借弹性、可扩展性和高可用性,提供便捷的计算环境;分布式系统架构则通过多计算机协同工作,实现任务并行和容错。两者相互依存,共同推动企业数字化转型、科技创新、公共服务升级及数字经济发展。虚拟化、分布式存储和计算、网络技术是其核心技术。未来,深化研究与应用这些技术将促进数字化时代的持续进步。
670 4
|
人工智能 分布式计算 DataWorks
分布式×多模态:当ODPS为AI装上“时空穿梭”引擎
本文深入探讨了多模态数据处理的技术挑战与解决方案,重点介绍了基于阿里云ODPS的多模态数据处理平台架构与实战经验。通过Object Table与MaxFrame的结合,实现了高效的非结构化数据管理与分布式计算,显著提升了AI模型训练效率,并在工业质检、多媒体理解等场景中展现出卓越性能。
|
存储 自然语言处理 搜索推荐
分布式搜索引擎ElasticSearch
Elasticsearch是一款强大的开源搜索引擎,用于快速搜索和数据分析。它在GitHub、电商搜索、百度搜索等场景中广泛应用。Elasticsearch是ELK(Elasticsearch、Logstash、Kibana)技术栈的核心,用于存储、搜索和分析数据。它基于Apache Lucene构建,提供分布式搜索能力。相比其他搜索引擎,如Solr,Elasticsearch更受欢迎。倒排索引是其高效搜索的关键,通过将词条与文档ID关联,实现快速模糊搜索,避免全表扫描。
995 123
|
机器学习/深度学习 分布式计算 数据挖掘
MaxFrame 性能评测:阿里云MaxCompute上的分布式Pandas引擎
MaxFrame是一款兼容Pandas API的分布式数据分析工具,基于MaxCompute平台,极大提升了大规模数据处理效率。其核心优势在于结合了Pandas的易用性和MaxCompute的分布式计算能力,无需学习新编程模型即可处理海量数据。性能测试显示,在涉及`groupby`和`merge`等复杂操作时,MaxFrame相比本地Pandas有显著性能提升,最高可达9倍。适用于大规模数据分析、数据清洗、预处理及机器学习特征工程等场景。尽管存在网络延迟和资源消耗等问题,MaxFrame仍是处理TB级甚至PB级数据的理想选择。
498 6

热门文章

最新文章