中企动力 > 商学院 > 大数据计算平台
  • ?

    从内部业务到外部赋能,详解阿里大数据计算平台的扩张野心

    Rosalba

    展开

    2018 年是阿里巴巴公司成立的第 19 个年头。在这过去的 19 年里,伴随着中国互联网的快速发展,阿里巴巴也无到有、从小到大,迅速成长为一家世界级的互联网巨头,创造了一个令世界瞩目的「中国奇迹」。

    而在这 19 年的时间内,公众对于阿里巴巴公司的认知也在悄然发生着变化。从早年间的 B2B 公司到后来的 C2C(淘宝)、B2C(天猫)的电商公司再到现在一个无所不包的阿里巴巴生态体系,「阿里巴巴到底是一家什么公司?」这个问题可以有多个回答的角度,比如,阿里巴巴是一家以用户需求为导向的互联网公司,再比如,阿里巴巴是一家「商业太过成功以至于掩盖了技术创新的公司」(阿里巴巴 CTO 张建锋语)。

    但如果从最微观的角度切入,阿里巴巴其实一家大数据公司。在阿里所有的产品里,流淌的着是各种各样的数据,比如天猫淘宝的电商数据、阿里云的企业业务数据、支付宝的支付数据等等,这些海量的数据组成了阿里巴巴各个产品线,而让这些数据转化为业务和产品,最终成为可以让普通用户享受到的服务,则离不开一个稳定可靠的大数据计算平台,这也是阿里巴巴计算平台所要承担的艰巨任务。

    公开资料显示,阿里巴巴计算平台支撑了整个阿里经济体 90% 以上的结构化/非结构化数据的存储、交换、管控,数据规模已超 EB 级别。在上周的云栖大会上,阿里巴巴副总裁、计算平台负责人周靖人博士及其团队像外界展示了阿里巴巴大数据智能计算引擎的核心技术能力,比如可以实现海量数据规模下的高性价的离线实时计算,以及实时+离线任务一体化研发能力等等,这一系列新的能力也让其具备了新一代计算引擎的诸多特点。

    更重要的是,不管是大数据引擎 MaxCompute 还是实时计算引擎 Blink,都是在阿里内部被业务一步步「锻炼」出来的产品,因此具有实战性、可用性的优势。另一方面,作为阿里巴巴大数据研发平台的 DataWorks,在经过 9年 内部发展、5年公共云、3年专有云的发展后,也成为阿里巴巴大数据赋能行业的重要技术输出口。

    MaxCompute 与 Blink,从在线业务到民生业务的数据引擎

    先来看看 MaxCompute。这是阿里巴巴自主研发的大数据计算平台,从 2010 年开始正式开始运行在阿里云飞天分布式操作系统智商,提供统一的计算引擎,支持 SQL、MR、迭代计算、图计算、流计算。

    在历经多次、不同规模的业务锤炼后,目前 MaxCompute 承载了阿里巴巴集团内部 99% 的数据存储及 95% 的计算能力。

    与此同时,MaxCompute的成长速度也非常惊人。去年10月的云栖大会上,MaxCompute与TPC委员会的benchmark适配,在业界领先的基于端到端的大数据分析领域应用级测试基准下,MaxCompute完成了全球首次基于公共云的bigbench大数据基准测试,数据规模拓展到100TB,性能达到7830QPM,成为首个突破7000分的数据引擎。

    2018年,该性能测试的结果再次提升超过2倍,达到18176.71QPM。这一系列成绩充分展现了 MaxCompute 作为一款中国自主研发的大数据引擎,已经具备了可以引领行业发展的能力。

    再来看看看看 Blink。Blink 是阿里巴巴基于 Apache Flink 开源流处理框架所开发的实时计算引擎,过去三年,阿里的实时计算团队针对其内部特定的业务场景,对 Flink 做了大量优化迭代,并命名为 Blink。

    实时计算场景在电商业务里非常普遍,比如电商促销的场景,如何让用户的需求在短暂的促销阶段被更多地刺激出来,就考验着电商平台的搜索和推荐,这就需要电商平台的数据能在最短的时间内实现模型更新,这就是实时计算最能发挥作用的应用场景。

    在历年的双十一的大考中,公众最关注的 GMV 大屏幕的背后技术就是 Blink 实时计算引擎,每一条交易信息都是一个数据,从数据写入数据开始,到被实时处理并最终显现到大屏幕,都要求数据计算的精确性、可用性以及低延时(延迟在亚秒级别)。而双十一全天的活动里,每秒几十万笔的交易和支付的实时聚合统计操作全部是由Blink计算完成,从而最大限度地保证了双十一的稳定运行。

    从上文可以看出,MaxCompute 和 Blink 分别对应了不同领域的计算需求,前者主要应对海量数据的离线计算,而后者,则在实时计算中扮演重要角色,两个计算相辅相成,成为阿里巴巴内部诸多产品的底层数据支持平台。

    2016 年,阿里云推出 ET 城市大脑项目,在杭州,阿里云希望将城市交通数据统一到一个「大脑」中,通过云端的海量、实时计算,实现对城市发展的数字化管理,这也是对 MaxCompute 和 Blink 计算引擎的新考验,如果说过去的数据计算是处理互联网的交易数据,那么当数据范围扩大到物理世界,MaxCompute 和 Blink 能否有效应对呢?

    答案也很乐观。在上周发布的杭州城市大脑 2.0 中,阿里云 ET 城市大脑相的管辖范围扩大了28倍,优化信号灯路口1300个,覆盖杭州四分之一路口,同时已接入了视频4500 路。这意味着,MaxCompute 和 Blink 不仅可以计算互联网数据,还完全可以承载一个城市的离线和实时计算需求。

    这样灵活、强大的数据计算能力,也正在成为驱动其他行业变革的新变量。

    DataWorks,一站式数据开发平台

    事实上,MaxCompute 和 Blink 实时计算都已经运行在阿里云平台,企业和开发者可以根据自身需求去购买相应的服务。而在此次云栖大会上,阿里巴巴计算平台的多位技术专家还分享了 DataWorks 的数据研发平台对于更多行业的数据赋能能力。

    首先,DataWorks 的可用性已经得到验证。作为一个在阿里内部「孕育」出来的数据研发平台,DataWorks 也被广泛应用到阿里集团、蚂蚁金服、菜鸟、优酷、高德等所有事业部的数据开发流程里,还通过阿里云的公共云平台和专有云平台被广泛应用到多个国家和地区。

    其次,DataWorks 的技术能力毋庸置疑。不完全统计,2017年,以 DataWorks 为主体的阿里云数加,获得了国际软博会金奖;2018年,DataWorks 名列国家大数据博览会十佳产品,荣获最佳案例实践奖。

    2018 年 3 月,咨询机构 Forrester 发布 Cloud Data Warehouse 第一季度榜单,DataWorks 携手 MaxCompute,与AWS,Microsoft Azure,Google Cloud 一众强手共同进入云数仓第一阵营,是唯一入选榜单的中国企业,也奠定了世界级大数据研发平台的地位。

    第三,在产品设计上,DataWorks 拥有完整的开发流程,实现了端到端的数据开发。DataWorks 将上文提及的 MaxCompute 离线计算能力和 Blink 实时计算能力封装为可用的接口,另外还将阿里巴巴机器学习平台 PAI 的机器学习能力融合到平台里,覆盖从数据计算到模型训练、线上数据服务,再到云上应用搭建的一站式云上大数据解决方案。

    另外,基于云上编程环境 Cloud IDW,DataWorks 还提供从 Sql、python,甚至 Java 的开发能力,这也意味着,开发者不必花费过多时间和精力去配置各种开发变量,只需将开发环境切换到云端,然后直接写代码就能快速搭建自己的产品。

    DataWorks 的上述能力也在体现在阿里巴巴计算平台日前举办的云上编程比赛中,各路选手需要利用DataWorks 快速搭建一个天气预报云端应用。

    第一步是离线数据导入和处理。选手们要将历史数据通过数据集成导入到MaxCompute 表,然后在 DataWorks 编写离线 SQL 进行数据预处理,处理后的数据在 PAI 机器学习平台通过引用内置的各种算法/模板进行建模、训练,并最终一键发布到EAS提供预测服务。

    第二步则是实时数据的接入和处理。将实时采集的气象数据通过数据集成导入到DataHub,然后在DataWorks编写实时SQL进行数据加工,加工后的实时数据和离线基础数据拖过简单拖拽就可以装载到Lightning引擎进行异构数据整合,并提供实时交互式查询服务。

    第三步构建应用。在DataWorks 的数据服务中,可快速的打通 EAS 服务和 Lightning 引擎并生成高性能的在线 API,同时在 AppStudio 中可无缝对接数据服务API;用可视化组件模板,简单几步配置就可以完成云上Web应用开发;另外AppStudio也提供了在线IDE环境可支持Java在线开发、编译、调试、运行、版本管理、多用户协同编辑等功能。

    尾巴:数据时代的红利

    无论承认与否,「数据是新时代的石油」已然成为行业共识,向数据要价值正在成为全社会各个行业的方法论。在这场数据智能的淘金热里,阿里将自己放在行业赋能者的位置,既有能提供处理海量数据的 MaxCompute,还有支撑双十一的实时计算引擎 Blink,也有面向机器智能开发的 PAI,而在这一系列产品的上层,也就是最接近企业、开发者的那一层,DataWorks 整合了所有的核心技术,并以友好的界面、一站式的流程展现给企业、开发者。

    如果阿里巴巴过去 19 年的努力,践行了「让天下没有难做的生意」的口号,那么,现在的阿里巴巴大数据计算平台上的这些产品,则正在努力实现「让天下没有计算不了的数据」的新愿景,这是阿里巴巴技术驱动型公司最直接的体现,也是数据时代企业、个人开发者的新红利。(完)

  • ?

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

    怜翠

    展开

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

    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、Linux系统安装

    一般使用开源版的Redhat系统--CentOS作为底层平台。为了提供稳定的硬件基础,在给硬盘做RAID和挂载数据存储节点的时,需要按情况配置。比如,可以选择给HDFS的namenode做RAID2以提高其稳定性,将数据存储与操作系统分别放置在不同硬盘上,以确保操作系统的正常运行。

    2、分布式计算平台/组件安装

    当前分布式系统的大多使用的是Hadoop系列开源系统。Hadoop的核心是HDFS,一个分布式的文件系统。在其基础上常用的组件有Yarn、Zookeeper、Hive、Hbase、Sqoop、Impala、ElasticSearch、Spark等。

    使用开源组件的优点:1)使用者众多,很多bug可以在网上找的答案(这往往是开发中最耗时的地方);2)开源组件一般免费,学习和维护相对方便;3)开源组件一般会持续更新;4)因为代码开源,如果出现bug可自由对源码作修改维护。

    常用的分布式数据数据仓库有Hive、Hbase。Hive可以用SQL查询,Hbase可以快速读取行。外部数据库导入导出需要用到Sqoop。Sqoop将数据从Oracle、MySQL等传统数据库导入Hive或Hbase。Zookeeper是提供数据同步服务, Impala是对hive的一个补充,可以实现高效的SQL查询

    3、数据导入

    前面提到,数据导入的工具是Sqoop。它可以将数据从文件或者传统数据库导入到分布式平台。

    4、数据分析

    数据分析一般包括两个阶段:数据预处理和数据建模分析。

    数据预处理是为后面的建模分析做准备,主要工作时从海量数据中提取可用特征,建立大宽表。这个过程可能会用到Hive SQL,Spark QL和Impala。

    数据建模分析是针对预处理提取的特征/数据建模,得到想要的结果。如前面所提到的,这一块最好用的是Spark。常用的机器学习算法,如朴素贝叶斯、逻辑回归、决策树、神经网络、TFIDF、协同过滤等,都已经在ML lib里面,调用比较方便。

    5、结果可视化及输出API

    可视化一般式对结果或部分原始数据做展示。一般有两种情况,行数据展示,和列查找展示。

    以上就简单介绍这么多,如果有小伙伴想了解和学习更多的大数据技术,可以私信小编索要资料

  • ?

    怎样通过大数据计算一个国家的命运?

    邹元柏

    展开

    这篇文章的题目或者可以叫《“XXXX”——大数据时代的一则寓言故事》,但因为“XXXX”是一个自造的成语,我对它所诠释的本篇文章要表达观点的准确性并不满意,所以就用了这么一个“疑问式的标题党句式”。另外,在看完文章后,希望你能发挥自己的想法在评论区将这个成语,造的更贴切!

    言归正传。当人们掌握了真实有效的数据,然后经过科学系统的归纳与总结,形成结论,最终指导实践,进而改善现实的弊端,使人类社会不断进化。这是科学的工作思路,正确的工作方法,是大数据时代里的一个我们都在追求并努力实践的高效的运行机制。

    既然大数据可以指导生活生产实践,进化人类社会(国家的宏观调控政策,供给侧改革),那么,通过它是不是也可以预测到潜藏在未来的某些客观的,无法改变的事实?也就是说:我们事先运用科学的数据统计与分析,断定事物的发展在某个时间必然会进入某个状态。其实现实生活中,这样聊天攀谈的句式比比皆是。

    抛出个大胆的论定:通过复杂的数据研究,一个国家将会在未来的哪一年灭亡?

    这似乎成了“大数据”本身的一个悖论,因为人们研究大数据的初衷就是用来指导实践,改善现实中的弊端,进化社会的。但现实告诉我们,糟糕的事情每天都在不可避免的发生。不管是一个什么样的国家,它总会有它过去的历史,有现在的各种客观存在,和冗杂的意识形态,有将来的规划和方向。但这“规划和方向”有时并不以人的意志为转移。

    从理论上,这就意味着:我抛出的大胆定论是一种潜藏的客观存在。

    做个大胆的假设:公元前206年,秦国灭亡是个必然的事实。

    这个假设本身又很吊诡,因为它早已经成为事实。用既定的事实来倒退导向它的原因是我们中学历史课本常见的解题方法,而那个正确的答案就是通过科学的分析“大数据”得来的。也就是说,我们凭什么认为秦朝是因为这些原因而灭亡的,凭的就是这些科学的“大数据”。

    历史课本中关于“秦朝灭亡的根本原因是什么?”这个问题的答案是:秦的暴政。(准确来讲应该是封建统治者对生产资料的大量占有及对劳动人民的残酷剥削。)

    当然,秦朝灭亡的原因有很多,所以,我们首先归纳它灭亡的根本原因,明确了根本原因,再开始围绕着它做数据统计。问题的答案是“秦的暴政”,那么,我们就要研究秦朝实行的是什么“政策制度”。

    秦制:

    1、秦统一文字,对于中央集权的统一、文化的传播和发展,都起到了重要作用;

    2、统一全国货币,又规定了统一的度、量、衡制度。促进了商品流通,促进发展;

    3、征发大批人力修缮长城,使人民苦于劳逸。

    4、下令史官烧掉记载秦国以外各国历史的史书,有敢私下讨论的人处死刑;

    ……

    历史影响:使人民之间的沟通、交流变得更加容易。但是以高压政治和残酷的刑法为主实行集权制度,又把人们的生活推向了水深火热之中。暴政剥削导致暴动反抗。

    这里所罗列的“秦制”是一个广义的定义,它的意思是“客观的存在于秦朝时期方方面面的制度”。如果秦施行的政策是“ABCD”,那“A”涉及到方方面面又是“1234”,而“1”又会更具体的反映出每个人在日常生活中的情绪和态度,甚至是影响到自然界的反应。至此,每个人的态度和情绪或自然反应就开始反作用于“1”,“1”开始导向“A”,“A”就直接影响到“制度本身”,合而为一,整个社会的一切都开始由量变走向质变。其实整个过程是一个无限“深和广”的程式,并非如我所列举的这么简单。

    好的量变导向好的质变,坏的量变导向坏的质变,然后好的质变和坏的质变开始进行综合较量(这个过程也是我们日常生活中最能直接感受到的变化,我们不习惯考虑深层的原因,只习惯在事物发展到最明显的阶段后做反应)。

    结果是:好的质变,败给了坏的质变,最终导致应有的客观的质变——秦朝的灭亡。

    倒退一下支撑“秦朝灭亡”这一结果的论据:

    通过对“ABCD”及“1234”等一些列客观存在的“大数据”的分析,我们认识了“秦朝灭亡”这件事情。那么,构成这组“大数据”的“小数据”是什么?只有掌握了方方面面的“小数据”,我们才能拿现实生活中的具体事件作参考。

    比如:在秦朝的某年某月某日,咸阳城内或任何一个地方,发生了一件可以影响“数据”走向的变量,可以是暴动事件,甚至是两个人之间的小摩擦,那么这件事情就是导向“秦朝灭亡”这一结果的一条数据。又或者秦朝皇帝是通过暴力还是禅让,或者其他的什么样方式登上帝位的?不同的登基方式会导向什么样的结果?这样,皇帝登基时的年龄,有利条件和无利条件是什么?在上朝过程中更多的商议了那方面的事情?等等,这些都是一条条可计算的数据。

    首先请原谅,接下来的小标题会让你本能的产生一种,被骗入观看植入广告的厌恶感。但是没有广告。

    如果有一款这样的APP,它能不能运算出一个国家将在何时灭亡?

    通过汇总世界各国不同历史时期的不同王朝的兴衰存亡的数据,更加清晰的未来预测将会呈现在我们面前。虽然这个数据要足够庞大,运算方式也足够复杂,但它们是数据,数据的命运在这样的一个“大数据结合便捷互联网统计和运算平台的时代”确实就是客观存在的,也是可以被逐渐认知的,更不可否认它是能足够支撑一个“结果发生”的。

    所以,这样智能化的工具被生产出来,就很荒诞。当政者在实行一条政策或者针对某件事作出某个反应动作后,这款智能化工具调出大数据库的数据,直接生成另一组导向结果的数据,点击“提交”,然后显示:某年某月某日,因何事,而灭亡。不知道这款工具能不能成为一鼎可以敲响的警钟。

    另外,文首提出的小请求,有好的替代,还请费心,它的反义词是“天方夜谭” 。

  • ?

    基于大数据平台的数据分析

    花间词

    展开

    标签 | 大数据 架构

    作者 | 张逸

    无论是采集数据,还是存储数据,都不是大数据平台的最终目标。失去数据处理环节,即使珍贵如金矿一般的数据也不过是一堆废铁而已。数据处理是大数据产业的核心路径,然后再加上最后一公里的数据可视化,整个链条就算彻底走通了。

    数据处理的分类

    如下图所示,我们可以从业务、技术与编程模型三个不同的视角对数据处理进行归类:

    业务角度的分类与具体的业务场景有关,但最终会制约技术的选型,尤其是数据存储的选型。例如,针对查询检索中的全文本搜索,ElasticSearch会是最佳的选择,而针对统计分析,则因为统计分析涉及到的运算,可能都是针对一列数据,例如针对销量进行求和运算,就是针对销量这一整列的数据,此时,选择列式存储结构可能更加适宜。

    在技术角度的分类中,严格地讲,SQL方式并不能分为单独的一类,它其实可以看做是对API的封装,通过SQL这种DSL来包装具体的处理技术,从而降低数据处理脚本的迁移成本。毕竟,多数企业内部的数据处理系统,在进入大数据时代之前,大多以SQL形式来访问存储的数据。大体上,SQL是针对MapReduce的包装,例如Hive、Impala或者Spark SQL。

    Streaming流处理可以实时地接收由上游源源不断传来的数据,然后以某个细小的时间窗口为单位对这个过程中的数据进行处理。消费的上游数据可以是通过网络传递过来的字节流、从HDFS读取的数据流,又或者是消息队列传来的消息流。通常,它对应的就是编程模型中的实时编程模型。

    机器学习与深度学习都属于深度分析的范畴。随着Google的AlphaGo以及TensorFlow框架的开源,深度学习变成了一门显学。我了解不多,这里就不露怯了。

    机器学习与常见的数据分析稍有不同,通常需要多个阶段经历多次迭代才能得到满意的结果。下图是深度分析的架构图:

    针对存储的数据,需要采集数据样本并进行特征提取,然后对样本数据进行训练,并得到数据模型。倘若该模型经过测试是满足需求的,则可以运用到数据分析场景中,否则需要调整算法与模型,再进行下一次的迭代。

    编程模型中的离线编程模型以Hadoop的MapReduce为代表,内存编程模型则以Spark为代表,实时编程模型则主要指的是流处理,当然也可能采用Lambda架构,在Batch Layer(即离线编程模型)与Speed Layer(实时编程模型)之间建立Serving Layer,利用空闲时间与空闲资源,又或者在写入数据的同时,对离线编程模型要处理的大数据进行预先计算(聚合),从而形成一种融合的视图存储在数据库中(如HBase),以便于快速查询或计算。

    场景驱动数据处理

    不同的业务场景(业务场景可能出现混合)需要的数据处理技术不尽相同,因而在一个大数据系统下可能需要多种技术(编程模型)的混合。

    场景1:某厂商的舆情分析

    某厂商在实施舆情分析时,根据基于需求,与数据处理有关的部分就包括:语义分析、全文本搜索与统计分析。通过网络爬虫抓取过来的数据会写入到Kafka,而消费端则通过Spark Streaming对数据进行去重去噪,之后交给SAS的ECC服务器进行文本的语义分析。分析后的数据会同时写入到HDFS(Parquet格式的文本)和ElasticSearch。同时,为了避免因为去重去噪算法的误差而导致部分有用数据被“误杀”,在MongoDB中还保存了一份全量数据。如下图所示:

    场景2:Airbnb的大数据平台

    Airbnb的大数据平台也根据业务场景提供了多种处理方式,整个平台的架构如下图所示:

    Panoramix(现更名为Caravel)为Airbnb提供数据探查功能,并对结果进行可视化,Airpal则是基于Web的查询执行工具,它们的底层都是通过Presto对HDFS执行数据查询。Spark集群则为Airbnb的工程师与数据科学家提供机器学习与流处理的平台。

    大数据平台的整体结构

    行文至此,整个大数据平台系列的讲解就快结束了。最后,我结合数据源、数据采集、数据存储与数据处理这四个环节给出了一个整体结构图,如下图所示:

    这幅图以查询检索场景、OLAP场景、统计分析场景与深度分析场景作为核心的四个场景,并以不同颜色标识不同的编程模型。从左到右,经历数据源、数据采集、数据存储和数据处理四个相对完整的阶段,可供大数据平台的整体参考。

  • ?

    全国首个区块链大数据加密计算平台UD数链开源项目主网上线

    尔槐

    展开

    “全国首个区块链大数据加密计算平台”UD数链开源项目主网近期在上海、杭州、北京、贵州等地的节点实现了多地联网,并启动成功。

    据悉,UD数链开源项目由大数据智能服务商富数科技在2017年底率先发起,在进行了原型开发和技术验证之后,将设计开源,随后获得了挖财、游族MobData、神州泰岳、前隆科技、数汇通、浙江大数据交易中心、品钛集团、嘉银金科、麦达数字等企业和机构的认可,并成为创始节点单位。

    UD数链开源项目首批创始节点单位

    UD数链开源项目发起方富数科技是专注于金融行业的大数据解决方案提供商,拥有行业领先的人工智能、大数据、区块链技术和产品研发能力。富数科技创始人兼CEO张伟奇曾在Intel实验室、Capital One等知名机构任职,张伟奇表示“大数据、人工智能、云计算等产业如今发展都非常迅速,尤其是大数据已经全面进入我们的生活,谁拥有大数据,谁就掌握了主动权。”

    在张伟奇看来,大数据行业的痛点表现在以下几个方面:数据确权与权益保护难;数据交易的可信和安全得不到保障;数据的质量和隐私保护等,而运用区块链技术能够解决这些困扰行业多年的痛点。

    UD数链负责人卞阳介绍道,UD数链开源项目是一个以安全多方计算为基础协议,采用区块链的共识机制构建对等协作网络,致力于打造一个可信、开放、互利、公允的大数据区块链生态系统。UD——United Data,即汇合数据的意思,数链可以让参与各方在保护自己的专有数据、保护用户隐私的条件下,消除数据壁垒,实现数据安全、合法地流通,充分发挥数据价值。

    UD数链开放网络

    提到UD数链与国内不少区块链项目的差异时,卞阳指出“UD数链是以跨公司的方式进行建设,每个贡献单位都仅有一个节点,不存在单一公司控制。”网络结构分散,合作相互对等。

    其次,UD数链依靠数据合约实现多方持续共赢的博弈。所谓数据合约,即链上两方或多方的一种可执行的约定;这种约定,不是一次性的,而是一种持续性的博弈。以数据共享的场景为例,数据合约可以采用达成共识的规则,通过质押通证、延迟兑付、持续统计数据价值、自动调整贡献和权益等多种手段,使数据共享的价值分配,在持续的链式记账中,趋于合理公平。

    第三,链上不存储敏感数据,但记录了可追溯的凭证。数据的确权、授权、请求、服务等等都有凭证,这些凭证记录,增加了不诚信方的造假成本,并且可用于事后审计。在合约内我们还建立了“挑战-证明”机制,可以最大程度上防止造假行为,从而使数据更加可信。

    第四,数据安全来自于安全的存储、计算、传输和安全的流程。UD数链建立了一整套的安全协议,以确保数据在各个环节的安全。而对于参与各方,则建立了在实名认证的体系下的匿名交易模式,保护各方的商业经营信息不被泄露。

    UD数链联合建模与安全多方计算

    据了解,在技术层面,UD数链联合浙江大数据交易中心等机构引入了几种合规技术方案,包括“脱敏加密SDK”和“用户授权凭证合约”等,使得链上数据的流通更加合法合规。

    在产品层面,即将上链的首批产品包括:“区块链行业数据”、“网贷逾期匿踪信息共享”、“用户授权凭证”、“风控评分数据”等。“首批产品上链之后,UD数链有望成为日交易量最高的区块链项目之一”卞阳说道。

  • ?

    “技术盲”都夸好的大数据算法平台

    穷街

    展开

    随着AI人工智能、人脸识别技术的兴起,大数据已经成为各行业、业务领域竞相探索的方向。大数据算法更是核心竞争力,直接把握着大数据应用的可实现性和精准度。形象一点,它像是一本武林秘籍,得之可得天下。

    顺应市场浪潮,大数据算法平台如雨后春笋般出现在人们眼前,但无一例外,算法平台有技术门栏,无IT基础的人员根本无法使用。例如,大数据建模、模型训练、算法调用等等,面对满屏的技术术语和编程语句,无基础人员只剩下摇头的份了…即使购买算法平台的机构或企业配备技术人员,但新的问题出现了,用户发现号称提供大数据算法的平台,分析维度仍停留在传统BI层级。

    值得一提的是,德塔经过不断努力探索新的产品设计和合理的算法使用模式,以解决场景问题为导向,封装算法,提供用户易于理解,方便使用的大数据算法平台。现已嵌入德塔的智慧中枢产品。

    非技术人员可以使用吗?Of Course!

    德塔的大数据算法平台,将虚无缥缈的大数据算法实体化,变为看得见,摸得到的算法模型。为了达到非技术人员也可以使用的目的,系统将编程、算法调用等有技术门栏的部分隐藏到后台程序。而呈现给用户的仅仅是简单的操作。我们的算法团队根据常见的业务场景,构建出丰富的算法模型库。用户只需要选择对应的算法模型,录入相关的数据,便能得到计算结果,“傻瓜式”操作,简直是“技术盲”的福音。

    算法的应用场景充足吗?满足不了我需求怎么办?

    不得不说,我们的算法团队通过深挖行业、调研等方式,了解用户业务需求,提供丰富多样的业务场景算法模型。但真正了解业务需求的永远只有用户本人,考虑到这个方面,德塔在算法平台中提供自定义算法模型的能力。同样采用“傻瓜式”操作,达到用户可以自定义算法模型的目的。我的算法模型我做主,满足不了需求就自定义一个!

    我是专业技术人员,给我来个专业的使用通道!

    大数据算法平台,怎么可能没有专业模式?!德塔的大数据算法平台,为专业人士开放了算法调用接口,通过注册,获取开发者权限,利用样例代码快速实现代码层级的算法调用,将复杂的算法构建逻辑交由平台来处理,让用户更加专注于业务领域的拓展。

    以上是德塔产品的现阶段的功能。为何说“现阶段”呢?因为德塔的产品理念是“挖掘业务数据价值,打造聪明的智慧应用!所以,德塔的大数据算法平台必将经过多次产品迭代,不断为用户提供更实用、更有价值、更易懂的功能,完善“数据智慧,触手可及”的“智慧+”生态!

  • ?

    大数据平台快速解决方案

    恋繁华

    展开

    内容来源:2017年5月13日,周末去哪儿架构师李锡铭在“Java开发者大会 | Java之美【上海站】”进行《大数据平台快速解决方案中》演讲分享。

    阅读字数:1891 | 4分钟阅读

    摘要

    大数据(big data),指无法在一定时间范围内用常规软件工具进行捕捉、管理和处理的数据集合,是需要新处理模式才能具有更强的决策力、洞察发现力和流程优化能力的海量、高增长率和多样化的信息资产。周末去哪儿架构师李锡铭根据自己的成功经验,为我们分享大数据平台快速解决方案。

    大咖演讲视频

    http://t/R9an7Rr

    搭建始末

    当时我们确定要做大数据的时候,有两种选型。第一种选型是用用原生的、开源的大数据技术,需要自己搭建;第二种是ODPS。

    后来我们选择了利用原生大数据,自己搭建一个大数据平台。因为我们已经有了一定的小积累,并且也想做一个大数据方面的技术沉淀。

    在移动互联网时代,用户所有的行为、浏览、记录和收藏等所有的数据,我们都会把它拿下来分析,前段时间阶段性沉淀的东西有多少,是对之前的一个总结。这个数据还能帮助我们进行深度挖掘,之后如何对不同用户分类,做一个精准化的营销定位。

    每个公司都会对这些数据进行报表级的展现。我们最开始的数据实现方式是把所有用户的行为数据放到传统的关系型数据库中,利用纯Java应用程序去读这张表。当计算某个指标的时候,还会关联若干张子表。这张主表大概有几千万,其它子表也是百万级甚至千万级的。如果单纯用Java去算的话,还要额外处理多线程。

    所以我们用传统的Java纯程序+关系型数据库去处理报表的时候,在存储和计算的性能上会出现问题,以至于报表需求越来越慢。

    在这样的大背景下,我们改成了使用大数据去处理这种场景。

    技术概览

    Hadoop是现在所有大数据计算存储的一个底层概念,后面所有衍生的大数据产品都是在Hadoop的基础上进行衍生的。

    这张图是目前大数据平台的架构。

    原生的Hadoop应该包含了Hdfs(文件存储)、Yarn(资源调度)和Mapreduce(算法)。

    Spark是类似于Mapreduce的一个计算框架,它在很多场景中的性能会比原生的Mapreduce好很多,尤其是迭代计算的时候,会有好几个数量级的提升。

    Sqoop是一个数据的迁移工具。

    Hive是对底层Hdfs系统的文件抽象出一个类似Mysql的关系型数据库,但大前提是它是在Hadoop这个大的语义下的关系型数据库。

    Oozie是一个任务编排和调度的框架。

    Hue是大数据的管理后台。

    Zookeeper是分布式协调工具。

    1.组件分类

    基础数据:Mysql,File。基础数据层是游离于大数据之外的概念,它是传统的数据来源。

    大数据存储:Hdfs、Hive。大数据存储是最基础的文件存储,在这基础上抽象出一个大数据的关系型数据库。

    大数据计算:Mapreduce、Spark、Sqoop。Mapreduce是原生的,Spark是新生的,Sqoop是数据转移的工具。

    大数据协调与调度:Yarn、Zookeeper、Oozie。Yarn是原生的,Zookeeper是一个分布式保证文件原子性的工具,Oozie是调度工具。

    大数据展现:Hue。Curd的展现层。

    2.典型执行流程

    最开始说过,我们遇到的问题是,Mysql的表存不下,计算也有问题。在这个场景下要把数据,从Mysql转到大数据,并利用大数据进行计算,最后做一个展现。

    它的流程是,首先通过Sqoop把Mysql的数据一次性或是增量的同步到一张Hive表里,用Hive Sql写好查询后,本质上Hive Sql会转化成Mapreduce任务再去执行,最后数据就展现出来了。

    很多时候后台的服务Control层会有入口和出口,我们需要把入口和出口的参数都记下来,方便以后排错或做统计方面的应用。

    在应用程序里,把这些消息定时写到消息队列中,用Spark定时读消息队列,并把这些读取到的消息按Spark的方式做一个编程。这个任务最终会被丢到Hadoop的底层计算里,然后用Yarn去调度,计算出结果,把这个结果写入Hive,这就完成了一次流式计算。

    3.Hue

    这里写了一个Hive Sql,与传统Mysql的写法几乎一样。Hive Sql写好以后点执行。它的过程是把Sql首先交给Hive去跑,Hive用自己的Sql解析引擎把这个任务翻译成Mapreduce,Mapreduce再用Yarn跑在Hadoop上,最终把结果跑出来。

    4.存储:Hadoop hdfs

    HadoopHdfs是基础的存储层。

    HadoopHdfs其实只包含了两种类型,一个是Namenode,一个是Datanode。Namenode是一个管理的节点,而datanode只负责数据的存储和冗余。

    5.计算:Mapreduce&spark

    Hadoop原生的计算框架是Mapreduce,而spark是一个新兴的计算框架,它更快更全面。

    6.资源管理器:yarn、Apache、hadoop yarn

    资源管理器的架构内包含rescource manager和node manager。Rescource manager是管理节点,node manager是work节点。

    把任务丢给rescource manager,它去把任务分发给每个节点,做一些状态的变换,最后把结果通过rescource manager汇总以后,处理完毕交给客户端。

    7.hive

    hive的架构并不是很复杂,上层是一些用户的API、web页面和命令行。它的核心是执行引擎,把sql翻译成大数据平台可以接受的任务。底层基于存储,它可以存在hdfs上。

    8.sqoop

    主要用于在hadoop与传统的数据库间进行数据的传递。

    9.ooize

    大数据任务编排调度。

    学习与使用路线

    如果想要学习一些大数据相关的东西,我推荐可以先掌握一些基础,然后找一个场景套进技术里,进行快速实践。在快速实践的过程中会发现很多问题需要解决,很多知识需要补充,所以要在实践中前行,在错误中补充。

    我的分享到此结束,谢谢大家!

  • ?

    什么是大数据和大数据平台?

    鲁靖荷

    展开

    “大数据”时下一个热门的词语,近几年来,关于大数据的著作和文章铺天盖地,似乎也在共同在传递一个信息:越来越多的行业、人士开始关注并实际探索大数据的应用,我们正在一起描绘着大数据巨大效用的蓝图,但在实践的路上,我们都孩子起步阶段小步前行。

    大数据根基于互联网,数据仓库、数据挖掘、云计算等互联网技术的发展为大数据应用奠定基础。对于任何一个大数据的从业者或初接触者,或者都会有个共同的感触:大数据很有用!但大数据是什么呢?

    今天就给大家讲解一下:

    对于大数据的定义,我们来引用3个比较差用的大数据定义:

    1)Gartner:需要信息处理模式才能具有更强的决策力,洞察发现力和流程优化能力的海量、高增长率很多样化的信息资产。

    2)IDC:海量的数据规模(Volunme)、快速的数据流转和数据体系(Velocity)、多样的数据类型(Variety)、巨大的数据价值(Value)。

    3)Wiki:或称巨量数据、海量数据、大资料,指所涉及的数据量规模巨大到无法通过人工,在合理时间内达到截取、管理、处理、并整理成为人类所能解读的信息。

    其他关于大数据的定义也大抵类型,我们可以用几个关键词对大数据做一个界定。

    首先,“大规模”,这种规模可以从两个维度来衡量,一是时间序列累积大量的数据,二是在深度上更加细化的数据。

    其次,“多样化”,可以是不同的数据格式,如文字、图片、视频等,可以是不同的数据类别,如入口数据,经济数据等,还可以有不同的数据来源,如互联网、传感器等。

    最后,“动态化”,数据是不停变化的,可以随着时间快速增加大量数据,也可以是在空间上不断移动变化的数据。

    这三 个关键词对大数据从形象上做了界定。

    但是还需要一个关键能力,就是“处理速度快”。如果这么大规模、多样化又动态变化的数据有了,但需要很长的时间去处理分析,那不叫大数据。从另一个角度,要实现这些数据快速处理,靠人工肯定是没办法实现的,因此,需要借助于机器实现。

    最终,我们借助机器,通过对这些数据进行快速的处理分析,获取想要的信息或者应用的整套体系,才能称为大数据。

    我们可以用下面的图示给大数据定义:

    这下,就知道什么是大数据了吧

  • ?

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

    苻秋玲

    展开

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

大数据计算平台

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

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

img

在线咨询

建站在线咨询

img

微信咨询

扫一扫添加
动力姐姐微信

img
img

TOP