中企动力 > 商学院 > 流式大数据计算
  • ?

    快速“吞吐”大数据——前瞻计算机“高通量”时代

    鱼啊鱼

    展开

    新华社杭州10月26日电 题:快速“吞吐”大数据——前瞻计算机“高通量”时代

    新华社记者 董瑞丰、朱涵

    大数据与日俱增,计算机的“运力”能否跟上步伐?

    25日至27日在浙江杭州举办的中国计算机大会上,专家们探讨一种名为“高通量计算”的新生力量,描绘了未来计算世界的一幅新图景。

    新技术:天生擅长“吞吐”大数据

    什么是高通量计算?专家解释,就是同等时间内处理更多数据。简单来说,高通量计算擅长“并行”和“不规则运算”,能更好地应对海量数据。

    人类产生的数据量正以惊人速度递增。专业报告预测,从2016年到2025年,全球数据量将出现10倍的增长。中科院计算所高通量中心主任、中科睿芯董事长范东睿说,未来的计算机如果不能瞬时“吞吐”数据,人类将淹没在数据之中。

    “高通量计算擅长让数据在‘流动’中被处理,有效避免了反复访问带来效率与能耗上的损失,天生适合‘吞吐’大数据。” 范东睿说。

    中国科学院计算技术研究所所长孙凝晖认为,在未来的计算世界,两条技术路线将融合发展:一条是以传统高性能计算为核心的智能超算路线,一条是以高通量计算为核心的流式大数据路线,面对不同场景、解决不同问题,共生互补。

    新理念:计算机从“求快”转向“求多”

    有这样一个场景:把整个城市的交通红绿灯与摄像头联网,实时监测各路段的车流情况,随时调整红绿灯的时长,可以让道路更加通畅。

    要实现这一目标,对计算机的性能要求不见得很高,但需要计算机瞬时处理海量数据。从社交媒体到网络购物,从安防监控到无人驾驶,类似的场景未来将比比皆是。

    “计算机的大量任务正从‘求快’转向‘求多’。”范东睿说,“求多”不能仅靠设备堆积,大数据处理核心引擎上相应需要做出结构性的调整。

    国家计算机网络应急协调中心教授级高工包秀国认为,5G和物联网结合的时代临近,全社会的数据量将急剧增长,同时面临更实时的数据处理需求。在这种情况下,格外需要功能强大的、价格低廉的数据处理设备。

    “越来越多数据汇集到云端,高通量计算机变得非常有用。”腾讯云副总裁王龙说。

    新挑战:建设“信息高铁”

    在孙凝晖看来,高性能超级计算机好比飞机,速度快,但机场“吞吐”能力相对低;云计算中心则类似于公路,“吞吐量”大,但速度相对较慢。他把高通量计算机比作信息领域的高铁,既可以大量“吞吐”,又能快速运行。

    但如何平衡“吞吐量”、计算力和功耗之间的矛盾?专家指出,要用“异构计算”去解决这一问题,这涉及到全新的芯片结构和系统软件,非一日之功。

    中科院计算所于2005年在国内最早提出高通量计算的概念,经过10多年的前沿技术研究,才完成“低延迟、高通量、高确定性”的计算机原型系统研制。

    目前,高通量计算机正在针对深度学习、高通量音视频处理、科学大数据处理、信息安全检测、生物信息处理、大规模图数据处理等典型场景开展示范应用。

    责任编辑: 程瑶

  • ?

    新一代流式计算框架在金融科技中的应用

    尔冬

    展开

    大数据时代,数据计算已经渗透到了各行各业。业务沉淀数据,数据计算产生新的业务价值,数据计算正不断地用这种方式推动业务向前发展。

    大数据的计算模式主要分为批量计算(batch computing)、流式计算(stream computing)等,分别适用于不同的大数据应用场景。对于先存储后计算,实时性要求不高,同时数据规模大、计算模型复杂的应用场景,更适合使用批量计算。对于无需先存储,可以直接进行数据计算,实时性要求严格,但单次计算涉及数据量相对较小的应用场景,流式计算具有明显优势。

    与传统批量计算相比,流式计算的特点主要有以下几方面:

    无边界:数据到达、处理和向后传递均是持续不断的。

    瞬时性和有限持久性:通常情况下,原始数据单遍扫描,只有计算结果和部分中间数据在有限时间内被保存和向后传递。

    价值的时间偏倚性:随着时间的流逝,数据中所蕴含的知识价值往往也在衰减,即流中数据项的重要程度是不同的,最近到达的数据往往比早先到达的数据更有价值。

    金融行业是一种典型的流式计算应用领域,涵盖了包括用户行为分析、实时营销、个性化推荐、实时风控、实时反欺诈等多个计算场景。以实时金融风控场景为例,需要流式计算系统实时分析海量的用户行为数据,根据既定的规则计算出相应的指标,并与风险模型进行匹配,第一时间判断风险等级、发现异常事件,并作出相应的风险控制措施,自动告警通知、改变业务流程。

    流式计算框架的技术选型

    目前主流的流式计算框架有Storm、Spark Streaming、Flink三种,其基本原理如下:

    Apache Storm:以单事件来处理数据流(所有记录一个接一个处理),延迟性低(毫秒级),但消息保障能力弱,消息传输可能重复但不会丢失。

    Storm在运行中可分为spout与bolt两个组件,Spout是Stream的消息产生源,Bolt类接收由Spout或者其他上游Bolt类发来的消息,对其进行处理,以实现业务逻辑。Storm运行的应用程序叫做拓扑(Topology),拓扑中定义了Spout和Bolt的组合关系。

    Apache Spark Streaming:属于Spark API的扩展,以设定的时间间隔(如几秒种)处理一段段的批处理作业(即“微批处理”)。这种框架的延迟性较高(秒级),但能够保证消息传输既不会丢失也不会重复。

    spark程序是使用一个spark应用实例一次性对一批历史数据进行处理,spark streaming是将持续不断输入的数据流转换成多个batch分片,使用一批spark应用实例进行处理。

    Apache Flink:针对流数据+批数据的计算框架。把批数据看作流数据的一种特例,延迟性较低(毫秒级),且能够保证消息传输不丢失不重复。

    点击

    Flink创造性的统一了流处理和批处理,作为流处理看待时输入数据流是无界的,而批处理被作为一种特殊的流处理,只是它的输入数据流被定义为有界的。Flink程序由Stream和Transformation这两个基本构建块组成,其中Stream是一个中间结果数据,而Transformation是一个操作,它对一个或多个输入Stream进行计算处理,输出一个或多个结果Stream。

    从上述指标对比可以看出,Flink作为流式处理框架中的新兴之秀,兼具了storm的低延迟和spark的高吞吐、高消息保障特性。随着框架成熟度和社区活跃度越来越高,目前Flink已在国内外各大企业中有很多成功应用,很多企业纷纷将原有基于storm、spark streaming框架的流式计算系统转投为基于Flink,Flink成为当下流式计算框架的首选。

    流式计算在金融行业的应用架构

    顶象技术目前已为多个银行及互联网金融客户设计并搭建了基于Flink流式计算框架的一站式实时计算平台,满足金融客户实时风控、实时反欺诈等多个场景下的计算需求。

    金融行业的数据具有来源广泛、实时性要求高、吞吐量大、计算模型复杂等特性,顶象技术基于Flink引擎,进行了大量的功能及可用性设计,设计了一套可支持风控、反欺诈,流式计算和同步计算的一站式实时计算平台,满足金融行业实时风控、实时反欺诈等多个场景下的计算需求。

    顶象一站式实时计算平台的概要架构及特点:

    多数据源支持:支持从多种类型的消息中间件中获取消息流,以及从关系数据库或NoSql数据库中抽取全量/增量数据,并提供服务接口供其他系统实时推送数据。保证了计算数据的完整性和多样性。

    统一的异步/同步计算服务:基于Flink的流式计算仅能满足客户的实时异步计算需求,而无法满足需要实时返回计算结果的计算场景。因此,顶象一站式实时计算平台包含了高性能的同步计算框架,可满足同步返回计算结果的应用场景。该框架与异步流式计算框架结合,提供统一的计算服务,完美覆盖所有实时计算需求。

    可视化的计算配置和实时分析:为提高系统的易用性,使没有编程经验的业务人员也可通过流式计算快速实现业务目标,平台提供了可视化计算规则配置引擎,可针对各类应用场景,定义相应的计算规则。如针对风控场景,可定义各类风控事件的指标计算规则。同时,平台提供实时在线分析功能,业务人员在线选择数据源,设置报表统计规则,自动生成并实时刷新各种维度的图表,达到了数据的实时可视化分析。

    目前,顶象技术的一站式实时计算平台已在多个知名银行及互联网金融企业中部署上线,结合顶象实时决策引擎,实现了毫秒级风控指标计算与风控预警,同时为行内其他系统提供了统一实时数据计算服务,达到了很好的应用效果。

    金融机构在应用大数据的过程中,不仅需要看到其对于海量数据的存储、查询和分析类场景的需要,更需要探索如何运用多种大数据技术范式为业务提供解决方案。在应用较多的领域如用户画像、精准营销、实时反欺诈等领域,实时流计算的运用已经成为事实标准,未来在量化交易、风险检测、实时机器学习、实时决策引擎、设备异常分析等领域,相信还有更多用武之地。

  • ?

    【独家】一文读懂大数据计算框架与平台

    花缘灭

    展开

    1.前言

    计算机的基本工作就是处理数据,包括磁盘文件中的数据,通过网络传输的数据流或数据包,数据库中的结构化数据等。随着互联网、物联网等技术得到越来越广泛的应用,数据规模不断增加,TB、PB量级成为常态,对数据的处理已无法由单台计算机完成,而只能由多台机器共同承担计算任务。而在分布式环境中进行大数据处理,除了与存储系统打交道外,还涉及计算任务的分工,计算负荷的分配,计算机之间的数据迁移等工作,并且要考虑计算机或网络发生故障时的数据安全,情况要复杂得多。

    举一个简单的例子,假设我们要从销售记录中统计各种商品销售额。在单机环境中,我们只需把销售记录扫描一遍,对各商品的销售额进行累加即可。如果销售记录存放在关系数据库中,则更省事,执行一个SQL语句就可以了。现在假定销售记录实在太多,需要设计出由多台计算机来统计销售额的方案。为保证计算的正确、可靠、高效及方便,这个方案需要考虑下列问题:

    如何为每台机器分配任务,是先按商品种类对销售记录分组,不同机器处理不同商品种类的销售记录,还是随机向各台机器分发一部分销售记录进行统计,最后把各台机器的统计结果按商品种类合并?上述两种方式都涉及数据的排序问题,应选择哪种排序算法?应该在哪台机器上执行排序过程?如何定义每台机器处理的数据从哪里来,处理结果到哪里去?数据是主动发送,还是接收方申请时才发送?如果是主动发送,接收方处理不过来怎么办?如果是申请时才发送,那发送方应该保存数据多久?会不会任务分配不均,有的机器很快就处理完了,有的机器一直忙着?甚至,闲着的机器需要等忙着的机器处理完后才能开始执行?如果增加一台机器,它能不能减轻其他机器的负荷,从而缩短任务执行时间?如果一台机器挂了,它没有完成的任务该交给谁?会不会遗漏统计或重复统计?统计过程中,机器之间如何协调,是否需要专门的一台机器指挥调度其他机器?如果这台机器挂了呢?(可选)如果销售记录在源源不断地增加,统计还没执行完新记录又来了,如何保证统计结果的准确性?能不能保证结果是实时更新的?再次统计时能不能避免大量重复计算?(可选)能不能让用户执行一句SQL就可以得到结果?

    上述问题中,除了第1个外,其余的都与具体任务无关,在其他分布式计算的场合也会遇到,而且解决起来都相当棘手。即使第1个问题中的分组、统计,在很多数据处理场合也会涉及,只是具体方式不同。如果能把这些问题的解决方案封装到一个计算框架中,则可大大简化这类应用程序的开发。

    2004年前后,Google先后发表三篇论文分别介绍分布式文件系统GFS、并行计算模型MapReduce、非关系数据存储系统BigTable,第一次提出了针对大数据分布式处理的可重用方案。在Google论文的启发下,Yahoo的工程师Doug Cutting和Mike Cafarella开发了Hadoop。在借鉴和改进Hadoop的基础上,又先后诞生了数十种应用于分布式环境的大数据计算框架。本文在参考业界惯例的基础上,对这些框架按下列标准分类:

    如果不涉及上面提出的第8、9两个问题,则属于批处理框架。批处理框架重点关心数据处理的吞吐量,又可分为非迭代式和迭代式两类,迭代式包括DAG(有向无环图)、图计算等模型。若针对第8个问题提出来应对方案,则分两种情况:如果重点关心处理的实时性,则属于流计算框架;如果侧重于避免重复计算,则属于增量计算框架。如果重点关注的是第9个问题,则属于交互式分析框架。

    本文下面分别讨论批处理、流计算、交互式分析三种类别的框架,然后简要介绍大数据计算框架的一些发展趋势。文章最后介绍这一领域的学习资料。

    图1.大数据计算框架全景图2.批处理框架

    2.1.Hadoop

    Hadoop最初主要包含分布式文件系统HDFS和计算框架MapReduce两部分,是从Nutch中独立出来的项目。在2.0版本中,又把资源管理和任务调度功能从MapReduce中剥离形成YARN,使其他框架也可以像MapReduce那样运行在Hadoop之上。与之前的分布式计算框架相比,Hadoop隐藏了很多繁琐的细节,如容错、负载均衡等,更便于使用。

    Hadoop也具有很强的横向扩展能力,可以很容易地把新计算机接入到集群中参与计算。在开源社区的支持下,Hadoop不断发展完善,并集成了众多优秀的产品如非关系数据库HBase、数据仓库Hive、数据处理工具Sqoop、机器学习算法库Mahout、一致性服务软件ZooKeeper、管理工具Ambari等,形成了相对完整的生态圈和分布式计算事实上的标准。

    图2.Hadoop生态圈(删减版)

    MapReduce可以理解为把一堆杂乱无章的数据按照某种特征归并起来,然后处理并得到最后的结果。基本处理步骤如下:

    把输入文件按照一定的标准分片,每个分片对应一个map任务。一般情况下,MapReduce和HDFS运行在同一组计算机上,也就是说,每台计算机同时承担存储和计算任务,因此分片通常不涉及计算机之间的数据复制。按照一定的规则把分片中的内容解析成键值对。通常选择一种预定义的规则即可。

    执行map任务,处理每个键值对,输出零个或多个键值对。

    MapReduce获取应用程序定义的分组方式,并按分组对map任务输出的键值对排序。默认每个键名一组。

    待所有节点都执行完上述步骤后,MapReduce启动Reduce任务。每个分组对应一个Reduce任务。

    执行reduce任务的进程通过网络获取指定组的所有键值对。

    把键名相同的值合并为列表。

    执行reduce任务,处理每个键对应的列表,输出结果。

    图3.MapReduce处理过程

    在上面的步骤中,应用程序主要负责设计map和reduce任务,其他工作均由框架负责。在定义map任务输出数据的方式时,键的选择至关重要,除了影响结果的正确性外,也决定数据如何分组、排序、传输,以及执行reduce任务的计算机如何分工。前面提到的商品销售统计的例子,可选择商品种类为键。MapReduce执行商品销售统计的过程大致如下:

    把销售记录分片,分配给多台机器。每条销售记录被解析成键值对,其中值为销售记录的内容,键可忽略。

    执行map任务,每条销售记录被转换为新的键值对,其中键为商品种类,值为该条记录中商品的销售额。

    MapReduce把map任务生成的数据按商品种类排序。

    待所有节点都完成排序后,MapReduce启动reduce任务。每个商品种类对应一个reduce任务。

    执行reduce任务的进程通过网络获取指定商品种类的各次销售额。

    MapReduce把同一种商品下的各次销售额合并到列表中。

    执行reduce任务,累加各次销售额,得到该种商品的总销售额。

    上面的过程还有优化的空间。在传输各种商品每次的销售额数据前,可先在map端对各种商品的销售额进行小计,由此可大大减少网络传输的负荷。MapReduce通过一个可选的combine任务支持该类型的优化。

    2.2.DAG模型

    现在假设我们的目标更进一步,希望知道销售得最好的前10种商品。我们可以分两个环节来计算:

    统计各种商品的销售额。通过MapReduce实现,这在前面已经讨论过。对商品种类按销售额排名。可以通过一个排序过程完成。假定商品种类非常多,需要通过多台计算机来加快计算速度的话,我们可以用另一个MapReduce过程来实现,其基本思路是把map和reduce分别当作小组赛和决赛,先计算各分片的前10名,汇总后再计算总排行榜的前10名。

    从上面的例子可以看出,通过多个MapReduce的组合,可以表达复杂的计算问题。不过,组合过程需要人工设计,比较麻烦。另外,每个阶段都需要所有的计算机同步,影响了执行效率。

    为克服上述问题,业界提出了DAG(有向无环图)计算模型,其核心思想是把任务在内部分解为若干存在先后顺序的子任务,由此可更灵活地表达各种复杂的依赖关系。Microsoft Dryad、Google FlumeJava、Apache Tez是最早出现的DAG模型。Dryad定义了串接、全连接、融合等若干简单的DAG模型,通过组合这些简单结构来描述复杂的任务,FlumeJava、Tez则通过组合若干MapReduce形成DAG任务。

    图4.MapReduce(左)与Tez(右)执行复杂任务时对比

    MapReduce的另一个不足之处是使用磁盘存储中间结果,严重影响了系统的性能,这在机器学习等需要迭代计算的场合更为明显。加州大学伯克利分校AMP实验室开发的Spark克服了上述问题。Spark对早期的DAG模型作了改进,提出了基于内存的分布式存储抽象模型RDD(Resilient Distributed Datasets,可恢复分布式数据集),把中间数据有选择地加载并驻留到内存中,减少磁盘IO开销。与Hadoop相比,Spark基于内存的运算要快100倍以上,基于磁盘的运算也要快10倍以上。

    图5.MapReduce与Spark中间结果保存方式对比

    Spark为RDD提供了丰富的操作方法,其中map、 filter、 flatMap、 sample、groupByKey、 reduceByKey、union、join、cogroup、mapValues、sort、partionBy用于执行数据转换,生成新的RDD,而count、collect、 reduce、lookup、save用于收集或输出计算结果。如前面统计商品销售额的例子,在Spark中只需要调用map和reduceByKey两个转换操作就可以实现,整个程序包括加载销售记录和保存统计结果在内也只需要寥寥几行代码,并且支持Java、Scala、Python、R等多种开发语言,比MapReduce编程要方便得多。下图说明reduceByKey的内部实现。

    图6.RDD reduceByKey内部实现

    RDD由于把数据存放在内存中而不是磁盘上,因此需要比Hadoop更多地考虑容错问题。分布式数据集的容错有两种方式:数据检查点和记录数据的更新。处理海量数据时,数据检查点操作成本很高,因此Spark默认选择记录更新的方式。不过如果更新粒度太细太多,记录更新成本也不低。因此,RDD只支持粗粒度转换,即只记录单个块上执行的单个操作,然后将创建RDD的一系列变换序列记录下来,类似于数据库中的日志。

    当RDD的部分分区数据丢失时,Spark根据之前记录的演变过程重新运算,恢复丢失的数据分区。Spark生态圈的另一项目Alluxio(原名Tachyon)也采用类似的思路,使数据写入速度比HDFS有数量级的提升。

    下面总结Spark对MapReduce的改进:

    MapReduce抽象层次低,需要手工编写代码完成;Spark基于RDD抽象,使数据处理逻辑的代码非常简短。MapReduce只提供了map和reduce两个操作,表达力欠缺;Spark提供了很多转换和动作,很多关系数据库中常见的操作如JOIN、GROUP BY已经在RDD中实现。MapReduce中,只有map和reduce两个阶段,复杂的计算需要大量的组合,并且由开发者自己定义组合方式;Spark中,RDD可以连续执行多个转换操作,如果这些操作对应的RDD分区不变的话,还可以放在同一个任务中执行。MapReduce处理逻辑隐藏在代码中,不直观;Spark代码不包含操作细节,逻辑更清晰。MapReduce中间结果放在HDFS中;Spark中间结果放在内存中,内存放不下时才写入本地磁盘而不是HDFS,这显著提高了性能,特别是在迭代式数据处理的场合。MapReduce中,reduce任务需要等待所有map任务完成后才可以开始;在Spark中,分区相同的转换构成流水线放到同一个任务中运行。3.流计算框架

    3.1.流计算概述

    在大数据时代,数据通常都是持续不断动态产生的。在很多场合,数据需要在非常短的时间内得到处理,并且还要考虑容错、拥塞控制等问题,避免数据遗漏或重复计算。流计算框架则是针对这一类问题的解决方案。流计算框架一般采用DAG(有向无环图)模型。图中的节点分为两类:一类是数据的输入节点,负责与外界交互而向系统提供数据;另一类是数据的计算节点,负责完成某种处理功能如过滤、累加、合并等。从外部系统不断传入的实时数据则流经这些节点,把它们串接起来。如果把数据流比作水的话,输入节点好比是喷头,源源不断地出水,计算节点则相当于水管的转接口。如下图所示。

    图7.流计算DAG模型示意图

    为提高并发性,每一个计算节点对应的数据处理功能被分配到多个任务(相同或不同计算机上的线程)。在设计DAG时,需要考虑如何把待处理的数据分发到下游计算节点对应的各个任务,这在实时计算中称为分组(Grouping)。最简单的方案是为每个任务复制一份,不过这样效率很低,更好的方式是每个任务处理数据的不同部分。随机分组能达到负载均衡的效果,应优先考虑。不过在执行累加、数据关联等操作时,需要保证同一属性的数据被固定分发到对应的任务,这时应采用定向分组。在某些情况下,还需要自定义分组方案。

    图8.流计算分组

    由于应用场合的广泛性,目前市面上已经有不少流计算平台,包括Google MillWheel、Twitter Heron和Apache项目Storm、Samza、S4、Flink、Apex、Gearpump。

    3.2.Storm及Trident

    在流计算框架中,目前人气最高,应用最广泛的要数Storm。这是由于Storm具有简单的编程模型,且支持Java、Ruby、Python等多种开发语言。Storm也具有良好的性能,在多节点集群上每秒可以处理上百万条消息。Storm在容错方面也设计得很优雅。下面介绍Storm确保消息可靠性的思路。

    在DAG模型中,确保消息可靠的难点在于,原始数据被当前的计算节点成功处理后,还不能被丢弃,因为它生成的数据仍然可能在后续的计算节点上处理失败,需要由该消...

  • ?

    Apache Kafka:大数据的实时处理时代

    解智宸

    展开

    在过去几年,对于 Apache Kafka 的使用范畴已经远不仅是分布式的消息系统:我们可以将每一次用户点击,每一个数据库更改,每一条日志的生成,都转化成实时的结构化数据流,更早的存储和分析它们,并从中获得价值。同时,越来越多的企业应用也开始从批处理数据平台向实时的流数据数据平台转移。本演讲将介绍最近 Apache Kafka 添加的一些系统架构,包括 Kafka Connect 和 Kafka Streams,并且描述一些如何使用它们的实际应用体验。

    流处理

    在流处理刚被提出来的时候,很多人认为流处理只能进行做近似的结果或者增量的计算,倘若你想保证其安全性,以 Lamda 架构为基础,利用流处理得到最现在的结果。但同时你需要采用 batch processing 等其他方式来保证其全局的安全性以正确性。

    在如此多年的研究结果下,在我看来,流处理并不一定是近似的,或者是仅仅以无法保证真确性为代价而提高速度的一种数据处理方式。相反,流处理应该是一个与全局计算、batch processing 稍微有点不同的计算模型。跟批量处理不同之处在于,批量处理将数据引向计算,而流处理将计算引向数据。这句话大概有点模糊,接下来,我举几个大家熟悉的计算模型例子。

    第一个计算模型例子—请求应答模型。

    请求应答模型是业务生活中最常用的模型例子。首先提交一个请求到服务方,而服务方可能是一个数据库、也可能是别的存储工具;然后进行等待…等待;最后得到一个回答。这便是一次请求、一次计算、一次回答。该模型非常简单、也极易操作,当你需要延展到多个机器上时,只要简单地增加客户端以及处理器即可成功。但是缺点在于,不能达到大的吞吐量,每提交一次请求,都需要等待时间来获得最终应答的结果。

    第二种常见的模型就是批量处理如上图所示。如果请求应答模型在谱系的一端,那么 typo 的另一端则认为是批量处理。当我积累数据数量足够多的时候,一次性提交任务到数据仓库,再进行等待,等待时间短则几秒钟、几分钟,长则几小时,最后才得到最终的结果—所有输入对应的所有输出。该批处理模型的好处在于能够提高其吞吐率,一次的请求和应答可以得出较多结果。但它的缺点是具有高延时性,比如某数据产生时间为上午 6 点钟,用户点击某网页,由于批处理模型,每 12 小时才会运行一次,那么它必须等到上午 6 点到下午 6 点的所有数据完整以后才会进行工作,那么运行结果可能是用户点击的 12 个小时之后。高延迟性是批处理自身带有的特性。

    那么什么是流处理呢? 在我看来,流处理就是介于请求应答和批处理之间的一种新型计算模型或者编程模型。流处理并不等待数据的完整性,或者说数据本没有完整性这一讲法,数据本身就是一个数据流,当每个数据流每产生一个新数据的时候立刻被计算出、进行返回,因此数据是源源不断地通向计算,并且源源不断有结果被输出。你可以设想,与等待数据完全完成之后发布到计算上相比,流处理就是将计算移到你数据发生地进行实时计算的方式。

    为什么很多人之前有这样一种错觉,他们认为流处理可能存在有丢包的情况、或者说只可以得到近似的结果,其实这是早期的一些数据流处理系统所自带的一些限制。因此以 Lamda 架构为基础,在流处理上需要讨论不同维度的取舍。接下里我将举三个例子,延迟、、成本和正确性。正如很多人之前提及的,在进行流处理时候,其大多数情况需要用时间来换取正确性,或者用更多的成本换取时间等等。

    第一个例子,说如果你需要做一个实时的 ETL 处理。而关于 ETL 处理不需要太小的延迟,为达到低成本的一种保证,我们可以忍受几分钟或者 1 分钟的延迟;但是,如果你正在进行一个实时的在线监测,存在着几毫秒的延迟,那么这时候可能更愿意选择花大量的金钱,或者采取一些可能不必要的 possibility 来达到一种低延迟的效果;第二个例子,假设你在做一个在线付费协议,它也是一个流处理平台。由于在线付费协议可能关乎到其机构,或者其公司的利益所在,因此你会说,我需要保证百分之百的正确性,我不希望有任何丢包情况;

    第三个例子,如果你是做一个实时的日志处理,实时收集所有日志,并将其导入 root,在这种情况下,你可能会说,为了降低成本,我愿意付出一小部分正确性的代价,即使不能达到 100%、达到 99.99%、达到 99.9%,这样的结果都可以接受。这本是用户在定义不同流处理应用或者业务的时候应该可以自己做出的选择。但比较遗憾的是,多数早期的流处理平台其实并没有给予用户该种选择,他们自身的设计理念,那就是为了低延迟直接放弃掉正确性,或者说为了更高的吞吐量直接放弃低延迟。

    以上是我想分享的关于流处理的一些误会认知,如果我的分享能够让大家带走两个答案的话,我希望这就是一个。我认为流处理仅仅是一种不一样的计算模型或者编程模型,它将计算带到数据上,而不是将数据引用到计算上,并且在流处理的时候,用户往往需要在正确性、延迟性、成本等不同的维度上做出选择。

    Kafka 的角色

    为什么当我们说到流处理的时候,很多人都在说 Kafka。大多数人在最早接触 Kafka 时会说,Kafka 就是一个分布式发布订阅的消息系统,但是如果我们去观察 Kafka 的最初一些设计特性可发现以下几点内容。第一点,它可以作为一个写在磁盘上的缓存来使用,或者说,并不是仅基于内存来存储流数据,它可以保证数据包不被及时消费时,依然可用且不被丢失;第二点,由于位移的存在提供了逻辑上的顺序,在同一个话题上,第一个数据比第二个数据最先被发布的时候,也可保证在消费时也是永远第一个数据比第二个数据先被消费;第三点,因为 Kafka 是一个公有的大数据中转站,就是说,所有的数据只要在 Kafka 上,永远可以在 Kafka 周围进行业务的开发或者认知事物的开发。接下来我将花费一些时间详细介绍这三点之间的关系。

    Kafka 不仅仅是一个订阅消息系统,同时也是一个大规模的流数据平台,那么它提供了什么呢?第一,提供订阅和发布消息;第二,提供一个缓存的流数据存储平台;第三,提供流数据的处理平台。今天,我将着重讨论流式计算在 Kafka 上面的应用。

    流式计算在 Kafka 上的应用主要有哪些选项呢?第一个选项就是 DIY,Kafka 提供了两个客户端 —— 一个简单的发布者和一个简单的消费者,我们可以使用这两个客户端进行简单的流处理操作。举个简单的例子,利用消息消费者来实时消费数据,每当得到新的消费数据时,可做一些计算的结果,再通过数据发布者发布到 Kafka 上,或者将它存储到第三方存储系统中。DIY 的流处理需要成本。打个比方,考虑数据的延迟性,考虑不同时间上的管理分配,正如很多人提到的 processing time,这将是我后文会重点提及的概念。以上这些都说明,利用 DIY 做流处理任务、或者做流处理业务的应用都不是非常简单的一件事情。

    第二个选项是进行开源、闭源的流处理平台。比如,spark。关于流处理平台的一个公有认知的表示是,如果你想进行流处理操作,首先拿出一个集群,且该集群包含所有必需内容,比如,如果你要用 spark,那么必须用 spark 的 runtime。因为他们划定了你作为一个流处理平台使用者需要用到的所有行为,比如,资源管理系统、参数调配系统、容器配置、代码封装、分发等,以上行为都已被该平台所限定。一旦你选择使用甲就必须用甲套餐装备,如果选择使用乙就必须使用乙套餐装备。有人不禁提出疑问,我能不能既选择流处理平台,又使用自己选择的,我能不能这样做呢?

    这个应用场景其实很普遍,举个例子,可异步式微服务处理。什么叫异步式微服务处理?假设 Kafka 作为一个缓存数据,在该缓存区含有很多不同的业务。打个比方,一个网店的机构可以有不同的组、不同的员工,有人负责销售、有人负责商品分发,有人负责价格管理、有人负责在线实时的限流监控,不同的组、不同的员工可能会以不同的时间,或者以不同的代码来更新他们的产品,只要拥有一个异步式缓存机制,即 Kafka,便可扩大该微服务,而不需要他们的任何一个组之间进行同步请求应答机制。

    在该微服务情况下,每个小组的喜好、特性并不一致,有的组表示我需要做流处理平台,从 Kafka 读数据,处理完再写回 Kafka,并且想要使用 EWS 把我的应用部署在云端大规模集群上;而另外小组表示我不需要那么复杂,我只是小规模数据,不希望起一个集群,只需起三个机器,并且每个机器有 1GB 内存足以,可进行手动控制操作,不需要资源管理器。那么我们能不能同时满足他们不同的需求呢? 答案就是我接下来要说的第三种选项。

    第三种选项是使用一个轻量级流处理的库,而不需要使用一个广泛、复杂的框架或者平台来满足他们不同的需求。在 Kafka 0.10 当中已发布轻量级流处理内容平台,我们可以设想,跟其他客户端发布者和消费者一样,它也是一个客户端,不同之处在于它是一个计算者客户端,一个好用的、功能强大的客户端,并且支持 state processing、Windows 延时的、异步的、甚至不同数据的调控。 最重要的是 Kafka 作为一个库,可以采用多种方法来发布流处理平台的使用。比如,你可以构建一个集群;你可以把它作为一个手提电脑来使用;甚至还可以在黑莓上运行 Kafka。以上都是尤其简单的运行库的概念。

    因此我们要做的事情与使用 Kafka 其他的客户端类似,比如发布者、消费者,只要在代码里边加入就可以使用各种各样的 API。当你要调配控制 Kafka Stream 应用的时候,选择最基础的 War File 来运行或者采用 Java、C,甚至资源管理器来运行都是可行的。因为 Kafka Stream 是一个轻量级流处理的库,可支持各种各样的运维方式。

    在我们看来,简单的就是美的,只有给用户提供最大的兼容性与最大的延展性,用户才能得到最好的用户体验。

    Kafka Stream 的编程语言

    如果接触过 Storm、Spark 等流处理平台的同学可以发现,它们与 Kafka Stream 高阶位 DSL 语言其实有相似之处。如上图所示,首先定义一个 Streams 流, Streams 是从 topic1 中的 topic 获取得到,即定义 Streams、处理 Streams、得到新的 Streams。比如,从 topic1 里面得到两个原始数据流,然后数据流进行 countByKey 得到新的数据流叫做 Counts。那么 counts.to(“topic2”) 是什么意思呢?在获取到新的数据流之后写回 Kafka topic2 内,启动 KafkaStreams 进程,与 Kafka producer、Kafka consumer 类似,让它来运行已定义计算。

    正如大家所了解的,API 的使用其实很简单。提供一个简单的 API,用户简单地写入运行逻辑即可运行。但是编程应用总是容易的,而它的复杂程度在于,一旦你开始运维该应用,当你想要把业务拓展到更大规模,或者业务出现变化,或者集群不稳定,需要强大的运维时,运维的程度便显得异常重要,最上面的编程可能只是冰山一角。Kafka Stream 的设计理念是最简单的就是最美的,包括 API、运维、debugging,以及各种各样的方式,都是希望给用户带来最简单的体验。它的核心思想就是把难问题直接给 Kafka 集群本身。

    Kafka 的介绍

    Kafka 的核心思想是什么?就是把这些消息全部存成一个有序日志,所有的消息发布者把消息发布到底端,从某一个逻辑上的位移开始顺序读取所有的消息。它的一个好处在于所有的读和写,尽管都是刷到磁盘上,但都是按照顺序进行,该方式对磁盘的使用比较有效,倘若消费者和发布者隔得比较近,将利用 page cash 直接读数据。

    延展性。如上图,提供 topic 以及 topic partitions,即话题与话题分区的机制。每个用户有不同的 topic,每个 topic 可以有多个分区,每个分区可被装载在不同的机器上,当用户提高规模之后,Kafka 只需要简单地增加机器和 topic partitions 数量,或者采用 ROM balance 的方式到不同机器上,即可达到线性延展方式。

    以上是 Kafka 最简单的核心思想,接下来我将介绍 Kafka Streams 作为 Kafka 客户端如何利用以上核心思想来设计流处理的平台。数据流其实就是有序的记录或消息,每个消息是一个 Key 加一个 Value,并且 record 与 Kafka 自身 massage 具有一一对应关系。

    用户所提供的业务上的计算模型,其实可用拓补结构进行表达。如上图,图的左边。用户首先进行定义数据流,然后对数据流进行计算,得到新的数据流,最终将数据流写回到 Kafka 内。每当用户进行定义的时候,每一步都会变成拓扑结构里面的一个点,每个点通过流进行计算,变成新的流来进行新的连接,最终在 Kafka 内部形成拓扑结构。用户并不需要在意该拓补结构,只需明白定义流、计算流、得到新的流,写回 Kafka。

    连接每一个不同的运算单元就是一个 Stream,即 record stream,每一个 Stream 都在源源不断地实时产生 record,每一个 record 是一个 key 加一个 value。利用 Stream Processor 连接 Stream,每个用户定义的流的一个计算单位对应着一个 Stream Processor。

    当用户定义每一步计算的时候,就是定义每个拓扑结构里面的每个点,最终把整个拓补结构定义完整到 Kafka Stream 来运行。计算单...

  • ?

    5个大数据处理/数据分析/分布式工具

    醉香

    展开

    0.Hadoop

    Hadoop是一个开源框架,它允许在整个集群使用简单编程模型计算机的分布式环境存储并处理大数据。它的目的是从单一的服务器到上千台机器的扩展,每一个台机都可以提供本地计算和存储。

    1.Druid

    Druid是实时数据分析存储系统,Java语言中最好的数据库连接池。Druid能够提供强大的监控和扩展功能。

    2.Ambari

    大数据平台搭建、监控利器;类似的还有CDH

    提供Hadoop集群

    Ambari为在任意数量的主机上安装Hadoop服务提供了一个逐步向导。Ambari处理集群Hadoop服务的配置。

    管理Hadoop集群

    Ambari为整个集群提供启动、停止和重新配置Hadoop服务的中央管理。

    监视Hadoop集群

    Ambari为监视Hadoop集群的健康状况和状态提供了一个仪表板。

    3.Spark

    大规模数据处理框架(可以应付企业中常见的三种数据处理场景:复杂的批量数据处理(batch data processing);基于历史数据的交互式查询;基于实时数据流的数据处理,Ceph:Linux分布式文件系统。

    4.Storm

    Storm是一个免费开源、分布式、高容错的实时计算系统。Storm令持续不断的流计算变得容易,弥补了Hadoop批处理所不能满足的实时要求。Storm经常用于在实时分析、在线机器学习、持续计算、分布式远程调用和ETL等领域。Storm的部署管理非常简单,而且,在同类的流式计算工具,Storm的性能也是非常出众的。

    以上5个大数据处理/数据分析/分布式工具,有任何IT问题都欢迎问我~这里有一些我收集的资料想要的评.论!回复【资料】即可

  • ?

    人人都在说的大数据到底是什么?(技术层)

    濮阳苡

    展开

    前沿技术普及系列·写在前面

    不知道大家有没有和象牙妹儿一样的感觉,便是最近1年很多AI产品铺天盖地的来到了我们的生活,比如某些家的智能陪伴音响、翻译棒、智能机器人、AR拍照手机等等。

    这里面背后的前沿技术大数据、云计算、区块链、物联网、 5G、数字化…前几年甚至到现在听着还很遥远,但却已有不少在冲击着商业的世界,渗透进我们的生活!

    如何拨开层层浓雾,抵达未来的彼岸?如何把握技术的钥匙,启开未来商业世界的大门!

    大象互联网圈发起主办的第二届中国【郑州】开发者大会不但关注技术·商业的融合,也关注着前沿技术的分享。

    因此,从今天起,大象互联网圈公众号将目光聚焦在大数据、云计算、物联网IOT、人工智能AI、区块链、VR/AR等板块,一一从概念层、目前行业发展状况、企业实际应用层、技术层、岗位前景等为大家带来干货的普及,快来和象牙儿妹一块围观吧!

    上一篇:人人都在说的大数据到底是什么?(概念层)

    本文由“壹伴编辑器”提供技术支持

    提到大数据技术,最基础和核心的仍是大数据的分析和计算。在2017年,大数据分析和计算技术仍旧在飞速的发展,无论老势力Hadoop还是当红小生Spark,亦或是人工智能,都在继续自己的发展和迭代。

    目前绝大部分传统数据计算和数据分析服务均是基于批量数据处理模型:使用ETL系统或OLTP系统进行构造数据存储,在线的数据服务通过构造SQL语言访问上述数据存储并取得分析结果。这套数据处理的方法伴随着关系型数据库在工业界的演进而被广泛采用。

    本文将分别讨论大数据技术涉及到的技术框架、平台以及未来的发展趋势。

    本文由“壹伴编辑器”提供技术支持

    1.批处理框架

    传统的批量数据处理模型通常基于如下处理模型:

    1.使用ETL系统或者OLTP系统构造原始的数据存储,以提供后续的数据服务进行数据分析和数据计算。用户装载数据,系统根据自己的存储和计算情况,对于装载的数据进行索引构建等一些列查询优化工作。

    因此,对于批量计算,数据一定需要加载到计算机系统,后续计算系统才在数据加载完成后方能进行计算。

    2.用户或系统主动发起一个计算作用并向上述数据系统进行请求。此时计算系统开始调度(启动)计算节点进行大量数据计算,该过程的计算量可能巨大,耗时长达数分钟乃至数小时。

    同时,由于数据累计的不可及时性,上述计算过程的数据一定是历史数据,无法保证数据的实时性。

    3.计算结果返回,计算作业完成后将数据以结果集形式返回用户,或者可能由于计算结果数量巨大保存着数据计算系统中,用户进行再次数据集成到其他系统。一旦数据结果巨大,整体的数据集成过程漫长,耗时可能长达数分钟乃至数小时。

    典型代表:Hadoop

    Hadoop是Apache的一个开源项目,是可以提供开源、可靠、可扩展的分布式计算工具。它主要包括HDFS和MapReduce两个组件,分别用于解决大数据的存储和计算。

    HDFS是独立的分布式文件系统,为MapReduce计算框架提供存储服务,具有较高的容错性和高可用性,基于块存储以流数据模式进行访问,数据节点之间项目备份。默认存储块大小为64M,用户也可以自定义大小。

    HDFS是基于主从结构的分布式文件系统,结构上包括NameNode目录管理、DataNode的数据存储和Client的访问客户端3部分。

    NameNode主要负责系统的命名空间、集群的配置管理以及存储块的复制;DataNode是分布式文件系统存储的基本单元;Client为分布式文件系统的应用程序。

    对于数据存储,HDFS采用的是多副本的方式来存储数据,即Client将数据首先通过NameNode获取数据将要存储在哪些DataNode上,之后这些存储到最新数据的DataNode将变更数据以同步或异步方式同步到其他DataNode上。

    在Hadoop3.0之后,采用Erasure Coding可以大大的降低数据存储空间的占用。对于冷数据,可以采用EC来保存,这样才能降低存储数据的花销,而需要时,还可以通过CPU计算来读取这些数。

    MapReduce是一种分布式计算框架,适用于离线大数据计算。采用函数式编程模式,利用Map和Reduce函数来实现复杂的并行计算,主要功能是对一个任务进行分解,以及对结果进行综合汇总。

    具体来说,MapReduce是将那些没有经过处理的海量数据进行数据分片,即分解成多个小数据集;每个Map并行地处理每一个数据集中的数据,然后将结果存储为,并把key值相同的数据进行归并发送到Reduce处理。

    本文由“壹伴编辑器”提供技术支持

    2.流计算框架

    不同于批量计算模型,流式计算更加强调计算数据流和低时延,流式计算数据处理模型如下:

    1.使用实时集成工具,将数据实时变化传输到流式数据存储(即消息队列,如RabbitMQ);此时数据的传输编程实时化,将长时间累积大量的数据平摊到每个时间点不停地小批量实时传输,因此数据集成的时延得以保证。

    2.数据计算环节在流式和批量处理模型差距更大,由于数据集成从累计变成实时,不同于批量计算等待数据集成全部就绪后才启动计算作业,流式计算作业是一种常驻计算服务,一旦启动将一直处于等待事件触发的状态,一旦小批量数据进入流式数据存储,流计算立刻计算并迅速得到结果。

    3.不同于批量计算结果数据需要等待数据计算结果完成后,批量将数据传输到在线系统;流式计算作业在每次小批量数据计算后可以立刻将数据写入在线系统,无需等待整个数据的计算结果,可以立刻将数据结果投递到在线系统,进一步做到实时计算结果的实时化展现。

    典型代表:Spark

    Spark是一个快速且通用的集群计算平台。它包含Spark Core、Spark SQL、Spark Streaming、MLlib以及Graphx组件。

    Spark SQL是处理结构化数据的库,它支持通过SQL查询数据。Spark Streming是实时数据流处理组件。MLlib是一个包含通用机器学习的包。GraphX是处理图的库,并进行图的并行计算一样。

    Spark提出了弹性分布式数据集的概念(Resilient Distributed Dataset),简称RDD,每个RDD都被分为多个分区,这些分区运行在集群的不同节点上。一般数据操作分为3个步骤:创建RDD、转换已有的RDD以及调用RDD操作进行求值。

    在Spark中,计算建模为有向无环图(DAG),其中每个顶点表示弹性分布式数据集(RDD),每个边表示RDD的操作。

    RDD是划分为各(内存中或者交换到磁盘上)分区的对象集合。在DAG上,从顶点A到顶点B的边缘E意味着RDD B是RDD A上执行操作E的结果。有两种操作:转换和动作。

    转换(例如;映射、过滤器、连接)对RDD执行操作并产生新的RDD。

    本文由“壹伴编辑器”提供技术支持

    3.交互式分析框架

    在解决了大数据的可靠存储和高效计算后,如何为数据分析人员提供便利日益受到关注,而最便利的分析方式莫过于交互式查询。

    这几年交互式分析技术发展迅速,目前这一领域知名的平台有十余个,包括Google开发的Dremel和PowerDrill,Facebook开发的Presto,Hadoop服务商Cloudera和HortonWorks分别开发的Impala和Stinger,以及Apache项目Hive、Drill、Tajo、Kylin、MRQL等。

    一些批处理和流计算平台如Spark和Flink也分别内置了交互式分析框架。由于SQL已被业界广泛接受,目前的交互式分析框架都支持用类似SQL的语言进行查询。早期的交互式分析平台建立在Hadoop的基础上,被称作SQL-on-Hadoop。

    后来的分析平台改用Spark、Storm等引擎,不过SQL-on-Hadoop的称呼还是沿用了下来。SQL-on-Hadoop也指为分布式数据存储提供SQL查询功能。

    典型代表:Hive

    ApacheHive是最早出现的架构在Hadoop基础之上的大规模数据仓库,由Facebook设计并开源。Hive的基本思想是,通过定义模式信息,把HDFS中的文件组织成类似传统数据库的存储系统。

    Hive保持着Hadoop所提供的可扩展性和灵活性。Hive支持熟悉的关系数据库概念,比如表、列和分区,包含对非结构化数据一定程度的SQL支持。它支持所有主要的原语类型(如整数、浮点数、字符串)和复杂类型(如字典、列表、结构)。

    它还支持使用类似SQL的声明性语言HiveQueryLanguage(HiveQL)表达的查询,任何熟悉SQL的人都很容易理解它。HiveQL被编译为MapReduce过程执行。下图说明如何通过MapReduce实现JOIN和GROUPBY。

    部分HiveQL操作的实现方式

    Hive与传统关系数据库对比如下:

    Hive的主要弱点是由于建立在MapReduce的基础上,性能受到限制。很多交互式分析平台基于对Hive的改进和扩展,包括Stinger、Presto、Kylin等。其中Kylin是中国团队提交到Apache上的项目,其与众不同的地方是提供多维分析(OLAP)能力。

    Kylin对多维分析可能用到的度量进行预计算,供查询时直接访问,由此提供快速查询和高并发能力。Kylin在eBay、百度、京东、网易、美团均有应用。

    本文由“壹伴编辑器”提供技术支持

    5.其他类型的框架

    除了上面介绍的几种类型的框架外,还有一些目前还不太热门但具有重要潜力的框架类型。图计算是DAG之外的另一种迭代式计算模型,它以图论为基础对现实世界建模和计算,擅长表达数据之间的关联性,适用于PageRank计算、社交网络分析、推荐系统及机器学习。这一类框架有GooglePregel、ApacheGiraph、ApacheHama、PowerGraph、,其中PowerGraph是这一领域目前最杰出的代表。很多图数据库也内置图计算框架。

    另一类是增量计算框架,探讨如何只对部分新增数据进行计算来极大提升计算过程的效率,可应用到数据增量或周期性更新的场合。这一类框架包括GooglePercolator、MicrosoftKineograph、阿里Galaxy等。

    另外还有像

    ApacheIgnite、ApacheGeode(GemFire的开源版本)这样的高性能事务处理框架。

    本文由“壹伴编辑器”提供技术支持

    6.总结与展望

    从Hadoop横空出世到现在10余年的时间中,大数据分布式计算技术得到了迅猛发展。不过由于历史尚短,这方面的技术远未成熟。各种框架都还在不断改进,并相互竞争。

    性能优化毫无疑问是大数据计算框架改进的重点方向之一。而性能的提高很大程度上取决于内存的有效利用。这包括前面提到的内存计算,现已在各种类型的框架中广泛采用。

    拥抱机器学习和人工智能也是大数据计算的潮流之一。Spark和Flink分别推出机器学习库SparkML和FlinkML。更多的平台在第三方大数据计算框架上提供机器学习,如Mahout、Oryx及一干Apache孵化项目SystemML、HiveMall、PredictionIO、SAMOA、MADLib。

    在同一平台上支持多种框架也是发展趋势之一,尤其对于那些开发实力较为雄厚的社区。

    Spark以批处理模型为核心,实现了交互式分析框架SparkSQL、流计算框架SparkStreaming(及正在实现的StructuredStreaming)、图计算框架GraphX、机器学习库SparkML。

    本文由“壹伴编辑器”提供技术支持

    7.学习资料

    最后介绍一下大数据计算方面的学习资料。

    论坛

    首推知乎、Quora、StackOverflow,运气好的话开发者亲自给你解答。其他值得关注的网站或论坛包括炼数成金、人大经济论坛、CSDN、博客园、云栖社区、360大数据、推酷、伯乐在线、小象学院等。

    微信订阅号

    InfoQ是最权威的,其他还有THU数据派、大数据杂谈、CSDN大数据、数据猿、Hadoop技术博文等,各人根据偏好取舍。

    官方网站文档

    若要进行系统的学习,则首先应参考官方网站文档。不少大数据平台的官方文档内容都比较详实,胜过多数教材。

    书籍

    国外O\'Reilly、Manning两家出版社在大数据领域出版了不少优秀书籍,特别是Manning的InAction系列和O\'Reilly的DefinitiveGuide系列。

    本篇关于大数据的技术层面分析就介绍到这,下一篇我们将对大数据的最终价值体现——大数据实践进行分析介绍。

    END

    第二届中国【郑州】开发者大会

    线上报名通道全面开放

    本次大会将开设从数字化到大数据落地专场,现已开放报名通道,门票分为49.9普通票和VIP票两种,其中VIP票权益非常丰富,票价为199元且仅限200张,下面是两种票的权益对比:

  • ?

    阿里专家强琦:流式计算的系统设计和实现

    静枫

    展开

    更多深度文章,请关注云计算频道:https://yq.aliyun/cloud

    阿里云数据事业部强琦为大家带来题为“流式计算的系统设计与实现”的演讲,本文主要从增量计算和流式计算开始谈起,然后讲解了与批量计算的区别,重点对典型系统技术概要进行了分析,包括Storm、Kinesis、MillWheel,接着介绍了核心技术、消息机制以及StreamSQL等,一起来了解下吧。

    增量计算和流式计算

    流式计算

    流计算对于时效性要求比较严格,实时计算就是对计算的时效性要求比较强。流计算是利用分布式的思想和方法,对海量“流”式数据进行实时处理的系统,它源自对海量数据“时效”价值上的挖掘诉求。

    那么,通常说的实时系统或者实时计算,严格意义上来说分成三大类:

    ad-hoc computing(数据的实时计算):计算不可枚举,计算在query时发生。

    stream computing(实时数据的计算):计算可枚举,计算在数据发生变化时发生。

    continuous computing(实时数据的实时计算):大数据集的在线复杂实时计算。

    增量计算

    增量计算是分批,也就是batch,每个batch会计算出一个function的delta值,数据的一个delta最终会变成对function的一个delta值,最终通过增量计算达到效果。

    batch => delta: f(x + delta) = g( f(x), delta )

    实际上是在数据的delta值上计算的一个结果,这个f(x)我们称之为oldValue,整个function的一个oldValue从公式就可以看到,整个增量计算与全量计算和批量计算有很大的不一样的地方,就在于它是有状态的计算,而批量计算系统和全量计算系统是无状态的计算,所以这就会导致整个系统的设计思路理念和整个的容错机制会有很大的不同,相对于oldValue本批次的数据,delta作为一个输入,整体上是一个有状态的计算,它会在系统的时效性、系统的复杂性和系统性能之间去做tradeoff,如果batch里的数据量是非常少的,那这个系统表现出来的时效性是最实时的,当然,整个系统的容错吞吐就会受到影响,就是说一批次的数据量是比较少的情况下,整个的系统吞吐会比较低,整个系统的容错复杂度也会比较高,那么在增量计算情况下,它有哪些优势呢?

    1.相比以前的全量计算,中间的计算结果是实时产出的,也就是说它的时效性是很强的;

    2.我们把一个计算平摊在每一个时间段,可以做到平摊计算。整个集群的规模是受峰值的影响,双十一的峰值流量是非常大的,如果按照最峰值的流量去计算,整个服务器资源是相对较高的,如果能够把传统的计算平摊在每一分钟每一秒,实际可以起到降低成本的作用;

    3.整个数据处理链路如果放在一次Query中进行处理,也即是全部的数据在进行一个function的计算时,会大量膨胀中间结果,也就是说像Group By Count会到达200G,而增量计算可以做到中间结果不膨胀;

    4.增量计算是一个有状态的计算,在分布式领域,有状态的failover策略会跟无状态的计算系统截然不同,但是它的优势是恢复快,任务可以切成很多碎片去运行,一旦任务因为任何几台服务器的抖动而宕机,整个的恢复是从前一次有效的batch开始计算,而不是像全量计算和离线计算一样,全部要重新进行计算,当在离线计算和在线计算混合部署的情况下,这显得尤为重要;

    5.增量计算把一大块数据分批去计算,因此在批量计算里面经常遇到会一些数据倾斜问题在增量计算并不会遇到。在真实场景下,数据倾斜会对整个计算系统产生非常致命的影响,所以假设不同的节点之间数据倾斜比是1000,这个实际是很平常的,双十一的时候,光小米一家店铺就做到了很高的销售额,小米店铺和其他店铺的成交是上万倍甚至几十万倍的scale,传统的分布式计算的整个计算延时是受最慢的那个节点影响,如果把全部的数据分批次,实际上对于每一批来说,数据的倾斜度就会缓解,而且每个批次是可以并行去运行的,所以这可以大大地去降低整个计算任务在数据倾斜情况下的运行效率问题。

    增量计算和流式计算应用场景

    日志采集和在线分析:如基于访问日志、交易数据的BI算法分析。比较有名的像Google的统计、百度的统计,一些网站根据访问日志,会分析出各种的UV、 PV、 IPV等运营指标,有了流式计算,就可以对这些访问的时效性做到秒级、分钟级的监控,比如双十一当天,不同的店铺会通过店铺的实时访问情况来决定后面的运营策略;

    大数据的预处理:数据清洗、字段补全等;

    风险监测与告警:如交易业务的虚假交易实时监测与分析;

    网站与移动应用分析统计:如双11运营、淘宝量子统计、CNZZ、友盟等各类统计业务;

    网络安全监测:如CDN的恶意攻击分析与检测;

    在线服务计量与计费管理系统 搜索引擎的关键词点击计费;

    此外,流式计算和增量计算也应用在工业4.0和物联网上。

    流式计算的数据特点

    流(stream)是由业务产生的有向无界的数据流。

    不可控性:你不知道数据的到达时机以及相关数据的顺序,对于数据质量和规模也是不可控的;

    时效性要求:在容错方案、体系架构和结构输出方面都与传统的计算是截然不同的;

    体系缺失:传统学术领域已经对批量计算和离线计算的体系研究的非常成熟,而在实时领域如数据仓库中间层等领域都是缺失的,包括数据源管理、数据质量管理等等。

    另外,数据处理粒度最小,可以小到几条数据,对架构产生决定性影响;

    处理算子对全局状态影响不同,有状态、无状态、顺序不同等;

    输出要求,比如一致性和连贯性等。

    整个流计算会对系统有非常多的不一样的要求,这就会导致整个系统有非常大的复杂性,跟离线非常的不同,我们的计算仍然要求时效性、要求快,质量上要求它的计算一定是精准的,对容错的要求,不论你的机器、集群、网络硬件有任何的宕机,计算应该是持续稳定,对整个计算的要求也是非常多样性的。关于多样性,不同的业务场景,对计算的结果要求也是不一样的,有些要求精确,一点数据都不能丢、精度损失,还有的业务场景要求可以多但是不能少,还有丢数据有一个sla在保证等,所以种种特点导致我们做流式计算和增量计算系统会面临与传统的离线计算和增量计算完全不同的要求。

    与批量计算的区别

    从架构角度,增量计算、流式计算和离线处理、批处理有什么本质的区别?

    与批量计算的区别如上图所示,比如全量计算设计理念是面向吞吐,而流式计算是实时计算的一部分,面向延时;随之而来的整个全量DAG是一个串型的DAG,是一个StageByStage的DAG,而流式计算的DAG是一个并行DAG,也就是说Batch跟Batch之间是完全可以并行的,离线的批量系统它的串型化和Streaming场景下的并行化,它们在整个数据的时效性上

    有非常大的区别,特别是在Latency的体现。

    典型系统计算概要分析

    下面将向大家介绍业界比较经典的几个流计算产品:

    Twitter Storm

    Storm是Twitter内部使用开源被广泛使用的一套流计算系统,那么它的一个核心概念是说,一个任务要创建一个Topology,它表示了一个完整的流计算作业,它的最开始的源头名字叫做Spout,做收集数据的任务,它的前面可以挂任何的数据源、任何一个队列系统甚至可以对接文件,那么Bolt是它的具体计算任务所在的载体,而Bolt里有诸多的Task,它是在Spout和Bolt里负责具体一个数据分片的实体,它也是Storm里调度的最小单位。Acker负责跟踪消息是否被处理的节点。Storm的整个容错是采用源头重发的消息机制

    源头重发在网络流量激增的情况下,会造成系统的雪崩风险大大提升。上图是两个Storm的作业,它先从源头读出数据,然后进行filter过滤,最终进行join,join后进行一些逻辑处理。

    Nimbus–Zookeeper–Supervisor

    Storm采用了Nimbus Supervisor之间的方式进行任务调度和跟踪,它们之间是利用Zookeeper来进行通讯,Nimbus相当于一个全局的任务Master,负责接收Topology,然后进行二重的资源调度,并且将调度的信息记录到Zookeeper中,定期检查Zookeeper中的各种Supervisor的心跳信息,根据心跳状态决定任务是否进行重新调度,而Supervsor充当着每台物理机的一个watchdog,它在轮询Zookeeper中的调度任务信息,然后接收到发现有启动任务的信息,就会拉启进程,启动Task,同时定期要把心跳信息写入Zookeeper,以便Supervisor来做出重新调度或者系统的重发操作。

    消息跟踪机制是Storm的核心,保证消息至少被处理一次,它追踪源头信息的所有子孙信息。

    基本思路如下:

    Acker节点是进行消息跟踪的节点,以源头消息的ID为hash key,来确定跟踪的Acker,源头信息对应的所有的子孙消息都有该Acker负责跟踪,而消息树上每产生一个新的子孙消息,则通知对应的Acker,子孙消息被处理,然后再去通知对应的Acker,当Acker里所有的子孙消息都被处理的时候,那么整个数据处理就完成了。

    子孙的产生是由父节点,而处理是被子节点。所以Storm用了一个非常巧妙的异或方法,当父节点产生这个消息时,产生一个随机数,把这个随机数异或到Acker里,Acker把这个随机数传递到下一步的节点,当这个节点正确被处理以后,再把这节点发送给Acker去做异或,所以Storm利用了这个Acker机制来压缩整个数据的跟踪机制,最终保证任意节点出现宕机而值不会变成0。

    Transactional Topology

    光有以上的机制,还远远不够。被系统重发的消息没有任何附加信息,用户无法判断消息是否是被重发的等一些问题还有待解决,为解决消息被重复处理的问题,Storm 0.7.0以后版本推出了Transactional Topology进行改进,

    原理如下:

    在Spout上将源头消息串行划分成 Batch,为每个Batch赋以递增的id,记录在Zookeeper中,利用Acker跟踪Batch是否被完全处理完成,超时或者节点异常,Spout重发Batch内的所有消息,不影响中间状态的操作可以并发的执行,例如 Batch内的聚合操作,用户代码利用唯一的Batch ID进行去重。

    整个Topology同一时刻只能有一个Batch正在提交,以保证在每个节点上Batch串行递增,简化用户去重的逻辑。

    Storm优缺点

    优点:消息在框架内不落地,处理非常高效,保证了消息至少被处理,Transactional Topology为消息去重提供了可能,调度模式简单,扩展能力强(关闭重发模式下),社区资源丰富,拥有各种常见消息源的Spout实现。

    当然Storm也有自己的劣势:Transactional Topology对Batch串行执行方式,性能下降严重;Batch太大太小都有问题,大小需要用户根据具体业务分情况设置等。

    Amazon Kinesis

    Kinesis系统是一种完全托管的实时处理大规模数据流的开放服务。

    所有节点运行于EC2中:相对Storm来说,它采用了消息节点内部重放的系统,而不是像Storm那样子源头重发,它的所有的节点都已经在EC2中,无需单独的调度策略、复用安全、资源隔离机制,且扩展性好、弹性可伸缩。

    只支持单级Task,可以利用多个Stream组成复杂的DAG-Task,用户代码需要实现DAG-Task内部的消息去重逻辑。

    数据收集与计算独立:数据收集模块(Shard)对消息进行持久化,最长保留24小时;可以Get方式从其它系统中读取Shard数据,计算模块(Kinesis App)处理被推送的数据,Instance个数与Shard个数相同;用户代码可以自主控制Checkpoint节奏。

    用户可以自主调用相应的SplitShard\MergeShard接口,Stream上所有App的并发度随之调整。具体实现方法如下:

    每个Shard串行将接收到的消息写入S3文件中,SplitShard后,原有Shard不再接收新数据,原有Shard对应的所有App的Instance处理完消息后关闭,启动新的Shards(两个)和对应新的Instances。

    使计算可以更加的弹性,服务的可用性也更高。

    Google MillWheel

    MillWheel系统是利用内部支持Snapshot功能的Bigtable来进行持续化中间结果,将每个节点的计算输出消息进行持久化,实现消息的“不丢不重”。

    区别于Storm的是,它没有复杂的跟踪树。因为每一级都把它的输出消息进行持久化,用户可以通过SetTimer\ProcessTimer接口解决用户代码在消息到来时才能取得控制流的弊端,然后在源头节点(Injector)上将数据打上系统时间戳,每个内部节点(Computation)计算出所有输入Pipe上的最小时间戳,向所有输出Pipe上广播当前完成的最小时间戳,用户可以利用Low Watermark这一机制解决消息乱序或一致性问题。

    核心技术

    那么,流式计算和增量计算中最核心的一些技术和难点有哪些呢?

    从这张图可以看到,整个流计算是由一个复杂的Topology所构成。那么,从输入到输出,其中比较重要的两个角色一是Jobmaster,一是Coordinator。Jobmaster是每个Job负责运行时的一个master;而Coordinator是刚才所说的消息跟踪的一个角色,所以Coordinator最好是完全可以做到无状态的线性扩展。

    Batch数据从源头进入后,进入Source节点,Source节点会从消息源读取数据,蓝色的部分代表着Worker节点,蓝色节点再向橙色节点进行数据传输的时候,遵循着Shuffle的方法,可以是哈希的方法,可以是广播的方法,也可以是任何用户自定义的方法,output节点会将输出结果向在线系统输出,或者向下一级MQ节点输出,输出的结果也是按照Batch去对齐。

    系统...

  • ?

    阿里新一代流式计算引擎 大数据培训Flink学习宝典奉上

    Haile

    展开

    5个月的好程序员大数据培训学习,只是冰山一角,对于大数据职业生涯,我们要走的路还很长。苦是真的,但是活着,身上的责任和梦想就应该去承担、去实现,要微笑的去面对磨砺。

    马上就要上战场了,今年毕业生820万,想想都可怕。付出不一定有结果,但是,不付出一定什么都没有!大数据学习内容杂而多,要系统的掌握整体,需要很多的时间。包括Apache官网的各个框架的熟悉,更是需要时间的沉淀。好在遇到了好程序员的负责讲师,整体课程安排也十分科学,以下是我对大数据Flink部分学习的一些总结:

    Flink是一个分布式流处理的开源框架,提供准确的结果,即使在无序或迟到数据的情况下也是如此,具有状态和容错能力,可以在保持一次性应用程序状态的同时无缝地从故障中恢复,大规模执行,在数千个节点上运行,具有非常好的吞吐量和延迟特性。

    此前,我们讨论了将数据集的类型(有界还是无界)与执行模型的类型(批量与流媒体)进行对齐。下面列出的许多Flink功能 - 状态管理,无序数据的处理,灵活的窗口 - 对于在无界数据集上计算精确的结果非常重要,并且由Flink的流式执行模型来实现。

    Flink保证有状态计算的exactly-once。“有状态的”意味着应用程序可以维护一段时间内已经处理的数据的汇总或汇总,并且Flink的检查点设置机制确保在发生故障时应用程序的状态exactly-once。Flink支持流处理和窗口事件时间semantics。事件时间可以轻松计算事件到达顺序不正确,事件可能延迟到达的流的精确结果。

    除了数据驱动的窗口,Flink还支持基于时间,计数或会话的灵活窗口。Windows可以通过灵活的触发条件进行定制,以支持复杂的流模式。Flink的窗口可以模拟数据创建环境的实际情况。

    Flink的容错功能是轻量级的,可以让系统保持高吞吐率,同时提供一次性一致性保证。Flink从零数据丢失的故障恢复,而可靠性和延迟之间的折衷可以忽略不计。

    Flink的保存点提供了一个状态版本管理机制,可以更新应用程序或重新处理历史数据,而且不会丢失状态,停机时间最短。

    Flink设计用于在数千个节点的大型集群上运行,除了独立集群模式之外,Flink还提供对YARN和Mesos的支持。

    希望我们能用大数据人工智能去改变这个世界!

  • ?

    这五种大数据计算框架,你一定要知道!

    安于

    展开

    随着这些年全世界数据的几何式增长,数据的存储和运算都将成为世界级的难题。之前小鸟给大家介绍过一些分布式文件系统,解决的是大数据存储的问题,今天小鸟给大家介绍一些分布式计算框架:

    Hadoop框架

    提起大数据,第一个想起的肯定是Hadoop,因为Hadoop是目前世界上应用最广泛的大数据工具,他凭借极高的容错率和极低的硬件价格,在大数据市场上风生水起。Hadoop还是第一个在开源社区上引发高度关注的批处理框架,他提出的Map和Reduce的计算模式简洁而优雅。迄今为止,Hadoop已经成为了一个广阔的生态圈,实现了大量算法和组件。由于Hadoop的计算任务需要在集群的多个节点上多次读写,因此在速度上会稍显劣势,但是其吞吐量也同样是其他框架所不能匹敌的。

    Storm框架

    与Hadoop的批处理模式不同,Storm采用的是流计算框架,由Twitter开源并且托管在GitHub上。与Hadoop类似的是,Storm也提出了两个计算角色,分别为Spout和Bolt。

    如果说Hadoop是水桶,只能一桶一桶的去井里扛,那么Storm就是水龙头,只要打开就可以源源不断的出水。Storm支持的语言也比较多,Java、Ruby、Python等语言都能很好的支持。由于Storm是流计算框架,因此使用的是内存,延迟上有极大的优势,但是Storm不会持久化数据。

    Samza框架

    Smaza也是一种流计算框架,但他目前只支持JVM语言,灵活度上略显不足,并且Samza必须和Kafka共同使用。但是响应的,其也继承了Kafka的低延时、分区、避免回压等优势。对于已经有Hadoop+Kafka工作环境的团队来说,Samza是一个不错的选择,并且Samza在多个团队使用的时候能体现良好的性能。

    Spark框架

    Spark属于前两种框架形式的集合体,是一种混合式的计算框架。它既有自带的实时流处理工具,也可以和Hadoop集成,代替其中的MapReduce,甚至Spark还可以单独拿出来部署集群,但是还得借助HDFS等分布式存储系统。Spark的强大之处在于其运算速度,与Storm类似,Spark也是基于内存的,并且在内存满负载的时候,硬盘也能运算,运算结果表示,Spark的速度大约为Hadoop的一百倍,并且其成本可能比Hadoop更低。但是Spark目前还没有像Hadoop哪有拥有上万级别的集群,因此现阶段的Spark和Hadoop搭配起来使用更加合适。

    Flink框架

    Flink也是一种混合式的计算框架,但是在设计初始,Fink的侧重点在于处理流式数据,这与Spark的设计初衷恰恰相反,而在市场需求的驱使下,两者都在朝着更多的兼容性发展。Flink目前不是很成熟,更多情况下Flink还是起到一个借鉴的作用。

    以上就是现在比较主流的大数据运算框架的介绍了,欢迎大家收藏转发。关注小鸟,获取更多大数据级相关技术的资讯与教程。

  • ?

    关于流式大数据实时处理技术、平台及应用

    粟米

    展开

    1 引言

    大数据技术的广泛应用使其成为引领众多行业技术进步、促进效益增长的关键支撑技术。根据数据处理的时效性,大数据处理系统可分为批式(batch)大数据和流式(streaming)大数据两类。其中,批式大数据又被称为历史大数据,流式大数据又被称为实时大数据。

    目前主流的大数据处理技术体系主要包括Hadoop[1]及其衍生系统。Hadoop技术体系实现并优化了MapReduce[2]框架。Hadoop技术体系主要由谷歌、推特、脸书等公司支持。自2006年首次发布以来, Hadoop技术体系已经从传统的“三驾马车”(HDFS[1]、MapReduce和HBase[3])发展成为包括60多个相关组件的庞大生态系统。在这一生态系统中,发展出了Tez、Spark Streaming[4]等用于处理流式数据的组件。其中,Spark Streaming是构建在Spark基础之上的流式大数据处理框架。与Tez相比,其具有吞吐量高、容错能力强等特点,同时支持多种数据输入源和输出格式。除了Spark开源流处理框架,目前应用较为广泛的流式大数据处理系统还有Storm[5]、Flink[6]等。这些开源的流处理框架已经被应用于部分时效性要求较高的领域,然而在面对各行各业实际而又差异化的需求时,这些开源技术存在着各自的瓶颈。

    在互联网/移动互联网、物联网等应用场景中,个性化服务、用户体验提升、智能分析、事中决策等复杂的业务需求对大数据处理技术提出了更高的要求。为了满足这些需求,大数据处理系统必须在毫秒级甚至微秒级的时间内返回处理结果。以国内最大的银行卡收单机构银联商务为例,其日交易量近亿笔,需对旗下540多万个商户进行实时风险监控,在确保这些商户合规开展收单业务的同时,最大限度地保障个人用户的合法权益。这样的高并发、大数据、高实时应用需求给大数据处理系统提出了严峻的挑战。银联商务以前使用的T+1事后风控系统存在风险侦测迟滞高(次日才能发现风险,损害已经造成)、处理时间长(十几个小时之后才能完成风险识别)、无法处理长周期历史数据(只能分析最近几日的流水数据)以及无法支持复杂规则(仅能支持累积求和等简单规则)等重大缺陷。为此,亟须研发全新的事中风控系统,以重点实现低迟滞(在1 min内甄别突发风险)、高实时(100 ms内返回处理结果)、长周期(可处理长达10年以上的历史周期数据)以及支持高复杂度规则(如方差、标准差、K阶中心矩、最大连续统计等)等目标。这一目标可以抽象为一个大数据处理科学问题:如何在一个完整的大数据集上,实现低迟滞、高实时的即席(Ad-Hoc)查询分析处理。

    2 技术解析

    现有的大数据处理系统可以分为两类:批处理大数据系统与流处理大数据系统。以Hadoop为代表的批处理大数据系统需先将数据汇聚成批,经批量预处理后加载至分析型数据仓库中,以进行高性能实时查询。这类系统虽然可对完整大数据集实现高效的即席查询,但无法查询到最新的实时数据,存在数据迟滞高等问题。相较于批处理大数据系统,以Spark Streaming、Storm、Flink为代表的流处理大数据系统将实时数据通过流处理,逐条加载至高性能内存数据库中进行查询。此类系统可以对最新实时数据实现高效预设分析处理模型的查询,数据迟滞低。然而受限于内存容量,系统需丢弃原始历史数据,无法在完整大数据集上支持Ad-Hoc查询分析处理。因此,研发具有快速、高效、智能且自主可控特点的流式大数据实时处理技术与平台是当务之急。

    实现一个融合批处理和流处理两类系统且对应用透明的系统级方案,需要攻克以下几个技术难点。

    (1)复杂指标的增量计算

    尽管计数、求和、平均等指标能够依靠查询结果合并实现,然而方差、标准差、熵等大部分复杂指标无法依靠简单合并完成查询结果的融合。再者,当查询涉及热点数据维度及长周期时间窗口的复杂指标时,多次重新计算会带来巨大的计算开销。

    (2)基于分布式内存的并行计算

    采用粗放的调度策略(例如约定在每天的固定时间将流数据导入批处理系统)会造成内存资源的极大浪费,亟须研究实现一种细粒度的基于进度实时感知的融合存储策略,以极大地优化和提升融合系统的内存使用效率。

    (3)多尺度时间窗口漂移的动态数据处理

    来自业务系统的数据查询请求会涉及多种尺度的时间窗口,如“最近5笔刷卡交易的金额”“最近10 min内密码重试次数”“过去10年的月均交易额”等。每次查询请求都重新计算结果会对系统性能造成极大的影响,亟须研究实现一种支持多种时间窗口尺度(数秒到数十年)、多种窗口漂移方式(数据驱动、系统时钟驱动)的动态数据实时处理方法,以快速响应来自业务系统的即席查询请求。

    (4)高可用、高可扩展的内存计算

    基于内存介质能够大大提升数据分析及处理能力,然而由于其易挥发的特性,一般需要采用多副本的方式来实现基于内存的高可用方案,这使得“如何确保不同副本的一致性”成为一个待解决的问题。此外,在集群内存不足或者部分节点失效时,“如何让集群在不间断提供服务的同时重新平衡”同样是一个待解决的技术难题。亟须研究分布式多副本一致性协议以及自平衡的智能分区算法,以进一步提升流处理集群的可用性以及可扩展性。

    “流立方”流式大数据实时处理技术在上述领域取得了一系列突破,该技术提供基于时间窗口漂移的动态数据快速处理,支持计数、求和、平均、最大、最小、方差、标准差、K阶中心矩、递增/递减、最大连续递增/递减、唯一性判别、采集、过滤等多种分布式统计计算模型,并且实现了复杂事件、上下文处理等实时分析处理模型集的高效管理技术。

    3 平台纵览

    基于“流立方”流式大数据实时处理技术,研发了“流立方”流式大数据实时处理平台。其应用框架如图1所示,具有良好的灵活性和适应性。平台的数据装载模块负责从具体业务系统中接入实时流数据,数据抽取模块负责批量抽取历史数据,模型装载模块负责将分析处理模型集中的计算模型和脚本加载到平台中。当收到业务系统发出的实时查询请求时,“流立方”平台能够根据分析处理模型在完整大数据集上实时计算出相应的指标,并进行判断,将结果反馈给业务系统。

    图1 “流立方”平台应用框架

    在测试环境为8台服务器(每台服务器配置24核 CPU、256 GB内存),同时计算16个统计指标(涉及4个维度,包含计数、求和、平衡、最大、最小、标准差、过滤、去重、排序、复杂事件处理等多种算法)的性能测试中,“流立方”平台达到了单节点写入大于43 000 TPS、8节点读取大于100万TPS、平均时延为1~2 ms的优异性能,如图2所示。

    图2 “流立方”平台性能指标

    “流立方”平台在解决批式大数据和流式大数据融合实时处理技术难题,实现优异性能的同时,还解决了流式大数据处理平台面临的两大工程化难题。一是作业的编排效率问题。大部分开源流处理平台在完成一个流处理编排时,都需要经过拓扑设计、代码编写、功能测试、打包部署等环节,一般需要一周的时间才能完成。“流立方”平台通过基于“所见即所得”的在线作业编排管理,将上线任务耗时降低到分钟级,大大提升了流处理作业的编排效率。二是流处理作业的灵活变更问题。流处理平台擅长进行逻辑预先定义的增量计算,尽管其计算效率极高,但计算灵活度受到限制。例如,某业务需要统计过去3个月的数据,现有的流处理平台在该业务上线3个月后才能完全生效,这样的工作方式使流处理技术在实际应用中受到很大的局限。“流立方”平台创新性地引入流媒体播放器的录制与重放思路,在原始数据进入流处理平台时,通过顺序写的方式持久化一份原始数据,在需要上线新的计算作业时,即刻重发指定时间窗口内的原始数据,从而实现快速(分钟级甚至秒级)计算作业上线。

    “流立方”平台引入了一系列创新技术,在性能、可用性、可扩展性等多个层面提升了流处理平台的处理能力,满足金融领域在内的众多领域的业务及运维需求。引入数据冲突智能规避技术,解决了流式处理中的热点数据处理问题,从而解决了大颗粒数据维度的处理效率问题;引入Paxos一致性协议,解决内存存储计算时多副本一致性问题,提供了面向运维人员透明的一致性解决方案;引入智能分区技术,基于一致性散列技术,进一步将散列值拆解为散列块,通过散列块的平滑迁移解决存储集群的可伸缩性设计问题,确保对于运维人员的集群变更透明性;引入计算作业的动态运行时加载技术,规避了作业手工打包部署的问题,进一步提升了开发人员的工作效率。

    在国内某大型银行卡收单机构组织的招标测试中,测试环节为两台低配置虚拟机,测试数据为该机构的数千万笔交易流水,计算逻辑包括50多条规则,涉及30多个统计指标。在该测试环节下,两家国外著名厂商中,一家厂商的计算时间长达24 h,另一家老牌数据库软件提供商则未能在一天内完成计算。相较于这些国外著名厂商的大数据处理平台,“流立方”平台能够在3 h内完成所有计算,且正确率为100%。

    4 应用场景

    “流立方”流式大数据实时处理系统在金融、交通、电信、公安等行业具有广泛的应用场景。以金融风控反欺诈为例,部署“流立方”风控系统仅需在交易前端增加风控探头,将实时交易数据旁路接入系统。“流立方”风控系统根据融合了专家知识和机器学习结果的数百条规则对每笔交易进行风险评估,判断是否允许进行该笔交易,流程如图3所示。该系统平均响应时间在6 ms以下,并发数超过50 000笔/s。同时,实现这一性能仅需要4台服务器。

    图3 基于“流立方”的金融风控反欺诈流程

    基于“流立方”的金融风控反欺诈技术体系包含技术(如设备指纹、代理侦测、生物识别、关联分析、机器学习等技术)、知识(如盗卡反欺诈、伪卡反欺诈、信用卡套现、营销反欺诈等规则与模型)、数据(如虚假手机数据、代理IP数据、P2P失信数据等标识数据)三大板块。技术部分中的设备指纹技术通过主被动混合的形式采集设备中软硬相关要素,结合概率论等算法为每一个设备颁发一个全球唯一的指纹编码,这些指纹编码在反欺诈的整个过程中起到非常积极的作用;代理侦测技术通过短时间内扫描IP相关端口来识别那些开启代理的IP,并在这些IP访问金融服务时进行识别;生物识别技术通过采集设备上用户的鼠标点击、触摸、键盘敲击等行为识别操作者是人还是机器以及是否操作者本人的问题;关联分析技术在底层通过图数据库存储不同节点以及关系信息,最终在界面上通过图的形式进行欺诈者关联分析及复杂网络分析;机器学习技术通过有监督、无监督的机器学习算法提升欺诈识别的准确率及覆盖率,并结合流立方技术提供模型的事中预测能力。

    基于上述技术体系,研发了银行业务风险实时监控系统、互联网支付业务风险实时监控系统、电商业务风险实时监控系统等金融风控反欺诈系列解决方案。这些方案已应用到银行、第三方支付机构、互联网金融等领域的上百家企业。目前50%以上的线下交易都在“流立方”的保护下进行,基于“流立方”的金融风控反欺诈解决方案每天为我国的金融机构抵御上亿次的攻击。该技术已经成为我国金融安全领域基础设施必不可少的组成部分。

    此外,在互联网机器防御系统中,“流立方”同样能发挥巨大作用。如今网络机器人遍布票务、电商、招聘、银行、政府、社交等各类网站,消耗了40%~60%的网络流量。网络机器人不仅消耗网络资源、影响正常客户访问、增加网站运营成本,还会爬取产品、价格信息,形成不正当竞争,甚至混淆网站用户生态,影响营销分析。传统的控制策略通过采取屏蔽频繁访问、设置验证码等方式防御网络机器人,无法应对日益智能化的新型网络机器人。基于“流立方”的互联网机器防御系统通过在Web服务器上嵌入插件或者独立的嗅探器(sniffer)程序,将全流量的Web访问请求旁路到独立的机器防御集群,进行实时的流量分析及防御决策,并将决策后的结果实时回馈到Web服务器插件中。Web服务器插件在判定当前访问的设备或者IP地址等是机器人时,能够自动改写响应内容,根据不同的风险级别自动拒绝交易或将访问者引导到第三方图形验证码服务商进行机器人验证。访问者在通过验证后可以继续正常访问Web服务。该系统还创新地将设备指纹以及人机识别服务运用到机器防御系统中,不仅增加了可分析维度,提升了控制颗粒度,同时能够对基于浏览器内核的高级爬虫进行防护。此外,将机器防御规则、数据服务、设备指纹、人机识别以及图形验证码以软件即服务(software as a service,SaaS)的形式提供服务,进一步降低了互联网网站客户的运维门槛,提升了产品竞争力。该机器防御系统工作过程如图4所示。

    基于“流立方”的实时机器防御系统通过多服务器访问流水关联决策、长周期数据决策、复杂规则爬虫识别、设备维度爬虫识别、人机识别等技术,实现了微秒级(400~800μs)的识别时延,同时具有机器人识别管控一体化、轻量级接入等优点。根据已经接入机器防御服务的几十家客户的反馈,基于“流立方...

流式大数据计算

所有视频需要登录后,才能观看

请先登录您的帐号,即可完整播放,如果您尚未注册帐号,请先点击注册。

img

在线咨询

建站在线咨询

img

微信咨询

扫一扫添加
动力姐姐微信

img
img

TOP