- ?
从大数据到快速数据:Kafka如何实现实时数据流
戚蓝
展开
转载自百家号作者:岛上IT
什么是Kafka?
Apache Kafka是一个开源项目,提供强大的连续数据流分布式处理 - 目前全球已有数千家企业信赖其产品,包括Netflix,Twitter,Spotify,Uber等。
技术,体系结构和实施使其高度可靠和高度可用,使流处理应用程序能够利用地理分布的数据流
实时数据流
通过实时数据,您可以实时处理并做出反应。
Kafka实时数据流可帮助您在事件发生时及时获得反馈,并将其用于您的优势。历史数据,加上来自Kafka的实时数据流,可帮助您为未来做出重要决策。实时数据可帮助您获得竞争优势,并使您可以更有效地使用大数据。有效的实时数据流是大多数现代应用程序的核心,架构设计不再只是一个数据队列。
像Instaclustr这样的托管服务提供商在其数据平台上提供Apache Kafka将提供全天候的专家支持,并持续监控吞吐量,延迟时间并全天候进行必要的调整,以便您继续享受Apache Kafka的优势
为什么超过三分之一的财富500强公司使用Kafka
Kafka同时支持写入和读取可扩展性。这意味着您可以将大量数据传输到Kafka,并同时执行消息的实时处理,包括将消息发送到其他系统。
这些应用程序实际上仅限于您的想象力。
哪些行业已经在使用Kafka?
kafka正在各个行业使用,包括物流,零售,医疗保健,金融服务,电子商务,物联网等。
例如在物流行业,Kafka正在帮助更快地移动软件包并帮助公司实现盈利。考虑到物流的真实世界复杂性,最好尝试跟踪货物,仓库和卡车的位置。当与这3个参数相关的实时数据通过Kafka管道时,可以收集有助于各种不同方面的信息,例如收集,存储,交付,计划和优化货物移动,实时检查,审计和欺诈检测。
同样,医疗保健行业中的患者医疗记录和医疗测试也是保险供应商的要求,以及设施管理,床位管理和患者EMR。Kafka管道有助于处理不同的情况。
- ?
足球比赛实时数据接口全平台api接入 ResonSports体育数据
被爱
展开
最近,雷速体育这个名字在体育行业内被越来越多的人提及。这支体育领域的新军在短短两年时间内迅速崛起,优秀的产品体验、令用户惊喜的产品功能及强大的技术支持能力,帮助雷速体育快速扩张着市场份额。
但是,在雷速体育的快速发展背后不为人知的是,ResonSports体育大数据依靠其自身掌握的核心数据处理技术,一直在为雷速体育提供全面的体育数据服务,尤其是面对雷速体育每日数千万级的体育数据吞吐量,ResonSports体育大数据依旧为雷速体育提供着稳定全面即时的体育数据服务。
经过两年和雷速在C端产品的沉淀,现在的ResonSports体育大数据,已掌握了成熟完善的体育数据支持技术,提供足球、篮球两大项赛事包括即时数据、历史统计分析、百家赔率、动画直播、情报优劣势分析等数据支持。
面对B端客户,现已实现所有足篮球数据的全平台api接口化,而且支持全套体育数据api接口的定制化服务,并且支持多种语言环境数据,可直接对接网站和手机app,方便您在用户端的数据快速部署,接口接入事宜,详询ResonSports客服。
而且在数据方面,ResonSports与欧洲多家数据服务商合作,拥有丰富的数据来源,经搭载云端传输的实时技术统计引擎对比整合处理,每场比赛择最优数据呈现。运营层面,ResonSports拥有十余人数据维护团队,全年365天7*24小时进行数据监控维护,保证数据的准确性与即时性。技术层面,ResonSports有着来自国内外各大厂资深工程师组成的技术团队,他们在体育数据处理方向有着非常丰富的经验。
凭借着兄弟品牌雷速体育在业内良好的口碑及强大的影响力,ResonSports已和百度体育、英超球队哈德斯菲尔德等多家机构达成合作关系。在体育产业崛起、体育数据产业还是一片蓝海的国内市场,他们相信会有一片属于ResonSports的广阔天地。
- ?
关于流式大数据实时处理技术、平台及应用
Galatea
展开
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)的识别时延,同时具有机器人识别管控一体化、轻量级接入等优点。根据已经接入机器防御服务的几十家客户的反馈,基于“流立方...
- ?
工业基础数据管理平台
梢雁
展开
数据传输(通讯网络)
数据传输类型分为有线传输、无线传输两种,有效的将数据传输到服务器的媒介。无线传输是以GPRS模块为主的专网和公网专线通讯,适用于现场复杂,距离较过,不宜布线和布线难度较大的环境。
工业基础数据管理平台的高度中心
工业基础数据管理平台包含数据存储系统、数据管理分析软件系统、调度台、大屏幕系统、视频监控系统、程控电话、机房及其配置等。
数据采集箱
根据用户提供测量点的位置分布确定采集箱站点的数量,按照采集点相对集中的原则,在全厂布置数据采集箱。
工业物联网网络配置图
数据管理软件功能介绍
为更好的做好工业基础数据管理平台,对项目中技术要求较高的数据查询分析系统软件部分,依托天津大学优秀的软件开发团队,在天津大学已开发的三维软件及电力管理软件的基础上,通过提出的要求,开发升级为配套的工业基础数据管理平台实用软件。该软件可以以厂区分布,行业工艺流程为基点,分析设置界面。软件具有灵活的柔性设置优势,用户可根据自身特点,建立形象的图片向导,建立自定义报表,建立数据需求,可直观体现能源使用分布情况及能源利用状况。使更多企业应用工业基础数据管理平台,达到节能降耗,提高能源使用效率,降低成本,同时掌控生产过程数据,实现产品质量的保证和稳定,并具有可追溯性,为中国智能工厂的建设做好基础工作。
预警显示
能源管理系统中预警显示,根据用户的需求对生产单位中重点的设备,重点的耗能介质进行预设置最大消耗量,如果超出该最大量,系统会自动弹出报警提示,并且有声音的提示,如果用户未处理该报警提示,则报警提示实时显示。
数据显示
按照设备查询、介质查询、单位查询或时间查询出厂区内所有数据点的表累计流量、瞬时流量等数值。
报表分析
按照客户需求生成能源报表,包括日报表、月报表、年报表、班报表等。报表内容为需求点信息,将生成的数据结果导出Excel。
图形分析
将数据通过图表,折线等图形形式显示,直观显示设备的能耗情况,还可以将现场不同设备按照任意时间和类型的历史数据进行图形显示。数据通过饼状图、柱状图、折线图、列表等可视化形式展示出来。
站点分布
能源管理系统可根据用户的需求对生产单位中各个站点进行检测,当有的站点发出报警信息时,报警站点的文字将会变成红色,便于即时进行处理。
耗能分析
根据用户实际用量,显示总的介质能耗,分析能耗弱项。
行业应用
电力需求侧管理平台
经济发展带来巨大的能源消耗,随着全球能源供应的日趋紧张,各种能源费用都呈上升趋势。近些年以电力为首的能源紧缺问题日益突出,如何科学、准确的提高电力能源的高度工作,确保电力能源供应质与量,同时减少电力能源供给问题对社会经济发展工作带来的压制早已成为了目前国内电力企业所要面临的一个重要问题。传统的手工抄表,费时、费力、准确性和及时性得不到可靠的保障,这就电力行业的企业管理不能获得足够详细和准确的原则数据。随着无线通信数字网络的发展,电力抄表能耗分析系统已成为发展的必然趋势。
电力抄表能耗分析系统采用B/S结构,支持智能手机等智能终端访问,能显示现场计量数据,所有数据在系统中可按时间进行数据历史查询,系统有能耗功能,可以按日、按周、按月查询各个单位中电量的能耗信息以及计算能耗消费价格,并可进行实际的调度,这样方便以后的各部门单位按计划进行实施,根据现实数据可发展并找出因人为因素所导致的能源消费,预测在采取了节能措施后,可能达到的预测效果。
冶金行业
冶金行业工艺流程
根据用户工艺流程,可设置多个工艺流程图,可分级设置工艺流程界面。软件具有灵活的柔性设置优势,用户可根据自身特点,建立形象图片向导,建立数据需求。
热力供热行业
随着国民经济的不断发展,人们对供暖质量的需求也在逐步提高。在传统供热模式下,为满足供热需求,换热站内设备运行参数多为人工调节,随着室外温度计热负荷的不断改变,不断的人工调节二次供水湿度以保证用户室内能够维持恒定的湿度,在大规模的应用与实验过程中,热风管控系统逐步走向了自动化的发展道路,且热网自动化控制的精细程度也在与日提升。
热力站智能热网管控系统针对每个热力站的运行情况进行远程数据记录、检测及控制调度。通过将每个热力站内控制箱接入本系统可对其实时数据,如供水温度、回水温度、多个阀门及水泵的状态等进行远程控制调节,使热力工作人员无需到每个热力站进行手动调节,达到方便快捷、提高效率、降低能源消耗的目的。
热电蒸汽管网管理平台
随着科学技术的日新月异,尤其是计算机、通讯技术的迅速发展,自动控制水平也得到了快速的发展和广泛的应用,需求用户对供热质量的要求不断提高和能源紧张的今天,提高供热质量同时节约能源势在必行。供热(蒸汽)管理能源管理系统将实时、全面了角供热系统的运行工况,监视不利于况点的运行,保证区域供热系统安全合理地运行,并可根据运行数据进行供热品质及舒适度,延长了设备的使用寿命。
供热(蒸汽)管理能源管理系统采用B/S结构,支持智能手机等智能终端访问,能显示能源介质当前采集数据,例如:瞬时流量、累积流量、温度、压力、电量等。所有数据可通过饼状图、柱状图、曲线图等可视化形式展示。系统具有超量报警功能,报警具有声音输出。所有数据存储到数据库内,支持导入和导出功能,具有数据筛选等辅助功能,数据的各项统计功能。具有数据曲线对比功能,支持报表输出,报表打印等功能。
这是鑫科推荐的第17个能源管理系统品牌
- ?
携程实时计算平台架构与实践丨DataPipeline
点绛唇
展开
文 | 潘国庆 携程大数据平台实时计算平台负责人
本文主要从携程大数据平台概况、架构设计及实现、在实现当中踩坑及填坑的过程、实时计算领域详细的应用场景,以及未来规划五个方面阐述携程实时计算平台架构与实践,希望对需要构建实时数据平台的公司和同学有所借鉴。
一、携程大数据平台之总体架构
携程大数据平台结构分为三层:
应用层:开发平台Zeus(分为调度系统、Datax数据传输系统、主数据系统、数据质量系统)、查询平台(ArtNova报表系统、Adhoc查询)、机器学习(基于tensorflow、spark等开源框架进行开发;GPU云平台基于K8S实现)、实时计算平台Muise;
中间层:基于开源的大数据基础架构,分为分布式存储和计算框架、实时计算框架;
离线主要是基于Hadoop、HDFS分布式存储、分布式离线计算基于Hive及Spark、KV存储基于HBase、Presto和Kylin用于Adhoc以及报表系统;
实时计算框架底层是基于Kafka封装的消息队列系统Hermes, Qmq是携程自研的消息队列, Qmq主要用于定单交易系统,确保百分之百不丢失数据而打造的消息队列。
底层:资源监控与运维监控,分为自动化运维系统、大数据框架设施监控、大数据业务监控。
二、架构设计与实现
1.Muise平台介绍
1)Muise是什么
Muise,取自希腊神话的文艺女神缪斯之名,是携程的实时数据分析和处理的平台;Muise平台底层基于消息队列和开源的实时处理系统JStorm、Spark Streaming和Flink,能够支持秒级,甚至是毫秒级延迟的流式数据处理。
2)Muise的功能
数据源:Hermes Kafka/Mysql、Qmq;
数据处理:提供Muise JStorm/Spark/FlinkCore API消费Hermes或Qmq数据,底层使用Jstorm、Spark或实时处理数据,并提供自己封装的API给用户使用。API对接了所有数据源系统,方便用户直接使用;
作业管理:Portal提供对于JStorm、Spark Streaming和Flink作业的管理,包含新建作业,上传jar包以及发布生产等功能;
监控和告警:使用Jstorm、Spark和Flink提供的Metrics框架,支持自定义的metrics;metrics信息中心化管理,接入Ops的监控和告警系统,提供全面的监控和告警支持,帮助用户在第一时间内监控到作业是否发生问题。
2.Muise平台现状
平台现状:
Jstorm 2.1.1、Spark 2.0.1、Flink1.6.0、Kafka 2.0;
集群规模:
13个集群、200+台机器150+Jstorm、50+Yarn、100+ Kafka;
作业规模:
11个业务线、350+Jstorm作业、120+SS/Flink作业;
消息规模:
Topic 1300+、增量 100T+ PD、Avg 200K TPS、Max 900K TPS;
消息延时:
Hermes 200ms以内、Storm 20ms以内;
消息处理成功率:
99.99%。
3.Muise平台演进之路
2015 Q2~2015 Q3 :基于Storm开发实时计算平台;
2016 Q1~2016 Q2 :Storm迁移JStorm、引入StreamCQL;
2017 Q1~2017 Q2 :Spark Streaming调研与接入;
2017 Q3~2018 Q1 :Flink调研与接入。
4.Muise平台架构
1)Muise平台架构
应用层:Muise Portal 目前主要支持了 Storm 与 Spark Streaming两类作业,支持新建作业、Jar包发布、作业运行与停止等一系列功能;
中间层:对底层Infrastructure做了封装,为用户提供基于Storm、Spark、Flink相对应的API以及各方面Services;
底层:Hermes & Qmq是数据源、Redis、HBase、HDFS、DB等作为外部的数据存储、Graphite、Grafana、ES主要用于监控。
2)Muise实时计算流程
Producer端:用户先申请Kafka的topic,然后将数据实时写到Kafka中;
Muise Portal端:用户基于我们提供的API做开发,开发完以后通过Muise Portal配置、上传和启动作业;作业启动后,jar包会分发到各个对应的集群消费Kafka数据;
存储端:数据在被消费之后可以写回QMQ或Kafka,也可以存储到外部系统Redis、HBase、HDFS/Hive、DB。
5.平台设计 ——易用性
首先:作为一个平台设计第一要点就是要简单易用,我们提供综合的Portal,便于用户自己新建管理它的作业,方便开发实时作业第一时间能够上线;
其次:我们封装了很多Core API,支持多套实时计算框架:
支持HermesKafka/MySQL 、QMQ;集成Jstorm、Spark Streaming、Flink;作业资源管控;提供DB、Redis、HBase和HDFS输出组件;基于内置Metric系统定制多项metric进行作业预警监控;用户可自定义Metric用于监控与预警;支持AtLeast Once 与Exactly Once语义。上文讲到平台设计要易用,下面讲平台的容错,确保数据一定不能出问题。
6.平台设计——容错
Jstorm:基于Acker机制确保At Least Once;
Spark Streaming:基于Checkpoint实现Exactly Once、基于Kafka Offset回溯实现At Least Once;
Flink:基于Flinktwo-phase commit + Kafka 0.11事务性支持实现Exactly Once。
7.Exactly Once
1)Direct Approach
当前大部分拿Spark Streaming消费Kafka的话,都是用Direct Approach的方式:
优点:记录每个批次消费的Offset,作业可通过offset回溯;
缺点:数据存储与offset存储异步:
数据保存成功,应用宕机,offset未保存 (导致数据重复);offset保存成功,应用宕机,数据保存失败 (导致数据丢失);2)CheckPoint
优点:默认记录每个批次的运行状态与源数据,宕机时可从cp目录恢复;
缺点:
1. 非100%保证ExactlyOnce;
https://iteblog/archives/1795 描述了无法保证Exactly once的场景;https://issues.apache.org/jira/browse/SPARK-17606 也存在doCheckPoint时出现块丢失的情况;2. 启用cp带来额外性能影响;
3. Streaming作业逻辑改变无法从cp恢复。
适用场景:比较适合有状态计算的场景;
使用方式:建议程序自己存储offset,当发生宕机时,如果spark代码逻辑没有发生改变,则根据checkpoint目录创建StreamingContext。如果发生改变,则根据实现自己存储的offset创建context并设立新的checkpoint点。
8.平台设计——监控与告警
如何能够第一时间帮用户发现作业问题,是一个重中之重。
集群监控
服务器监控:考量的指标有Memory、CPU、Disk IO、Net IO;平台监控:Ganglia;作业监控
基于实时计算框架原生Metric系统;定制Metrics反应作业状态;采集原生与定制Metrics用于监控和告警;存储:Graphite展 现:Grafana 告警:Appmon;我们现在定制的很多Metrics当中比较通用的是:
Fail:定期时间内,Jstorm数据处理失败数量、Spark task Fail数量;Ack:定期时间内,处理的数据量;Lag:定期时间内,数据产生与被消费的中间延迟(kafka 2.0基于自带bornTime)。携程开发了自己告警系统,将Metrics代入系统之后基于规则做告警。通过作业监控看板完成相关指标的监控和查看,我们会把Flink作为比较关心的Metrics指标,全都导入到Graphite数据库里面,然后基于前端Grafana做展现。通过作业监控看板,我们能够直接看到Kafka to Flink Delay(Lag),相当于数据从产生到被Flink作业消费,中间延迟是62毫秒,速度相对比较快的。其次我们监控了每次从Kafka中获取数据的速度。因为从Kafka获取数据是基于一小块一小块去获取,我们设置的是每次拉2兆的数据量。通过作业监控看板可以监控到每次从Kafka拉取数据时候的平均延迟是25毫秒,Max是 760毫秒。
接下来讲讲我们在这几年踩到的一些坑以及如何填坑的。
三、踩坑与填坑
坑1:HermesUBT数据量大,埋点信息众多,服务端与客户端均承受巨大压力;
解决方案:提供统一分流作业,基于特定规则与配置将数据分流至不同topic。
坑2:Kafka无法保证全局有序;
解决方案:如果在强制全局有序的场景下,使用单Partition;如果在部分有序的情况下,可基于某个字段作Hash,保证Partition内部有序。
坑3:Kafka无法根据时间精确回溯到某时间段的数据;
解决方案:平台提供过滤功能,过滤时间早于设定时间的数据(kafka 0.10之后每条数据都带有自己的时间戳,所以这个问题在升级kafka之后自然而然的就解决了)。
坑4:最初,携程所有的Spark Streaming、Flink作业都是跑在主机群上面的,是一个大Hadoop集群,目前是几千台规模,离线和实时是混布的,一旦一个大的离线作业上来时,会对实时作业有影响;其次是Hadoop集群经常会做一些升级改造,所以可能会重启Name Node或者Node Manager,这会导致作业有时会挂掉;
解决方案:我们采用分开部署,单独搭建实时集群,独立运行实时作业。离线归离线,实时归实时的,实时集群单独跑Spark Streaming跟Yarn的作业,离线专门跑离线的作业。
当分开部署后,会遇到新的问题,部分实时作业需要去一些离线作业做一些Join或 Feature的操作,所以也是需要访问主机群数据。这相当于有一个跨集群访问的问题。
坑5:Hadoop实时集群跨集群访问主机群;
解决方案:Hdfs-site.xml配置ns-prod、ns双重namespace,分别指向本地与主机群;
Spark配置spark.yarn.access.namenodes or hadoopFlieSystems
坑6:无论是Jstorm还是接Storm都会遇到一个CPU抢占的问题,当你上了一个大的作业,尤其是那种消耗CPU特别厉害的,可能我给它分开了一个Worker,一个CPU Core,但是它最后有可能会给我用到3个甚至4个;
解决方案:启用cgroup限制cpu使用率。
四、应用场景
1.实时报表统计
实时报表统计与展现也是Spark Streaming使用较多的一个场景,数据可以基于Process Time统计,也可以基于Event Time统计。由于本身Spark Streaming不同批次的job可以视为一个个的滚动窗口,某个独立的窗口中包含了多个时间段的数据,这使得使用SparkStreaming基于Event Time统计时存在一定的限制。一般较为常用的方式是统计每个批次中不同时间维度的累积值并导入到外部系统,如ES;然后在报表展现的时基于时间做二次聚合获得完整的累加值最终求得聚合值。下图展示了携程IBU基于Spark Streaming实现的实时看板。
2.实时数仓
1)Spark Streaming近实时存储数据
如今市面上有形形色色的工具可以从Kafka实时消费数据并进行过滤清洗最终落地到对应的存储系统,如:Camus、Flume等。相比较于此类产品,Spark Streaming的优势首先在于可以支持更为复杂的处理逻辑,其次基于Yarn系统的资源调度使得Spark Streaming的资源配置更加灵活,用户采用Spark Streaming实时把数据写到HDFS或者写到Hive里面去。
2)基于各种规则作数据质量检测
基于Spark Streaming,自定义metric功能对数据的数据量、字段数、数据格式与重复数据进行了数据质量校验与监控。
3)基于自定义metric实时预警
基于我们封装提供的Metric注册系统确定一些规则,然后每个批次基于这些规则做一个校验,返回一个结果。这个结果会基于Metric sink吐出来,吐出来基于metrics的结果做一个监控。当前我们采用Flink加载TensorFlow模型实时做预测。基本时效性是数据一旦到达两秒钟之内就能够把告警信息告出来,给用户非常好的体验。
五、未来规划
1.Flink on K8S
在携程内部有一些不同的计算框架,有实时计算的,有机器学习的,还有离线计算的,所以需要一个统一的底层框架来进行管理,因此在未来将Flink迁移到了K8S上,进行统一的资源管控。
2.Muise平台接入Flink SQL
Muise平台虽然接入了Flink,但是用户还是得手写代码,我们开发了一个实时特征平台,用户只需要写SQL,即基于Flink的SQL就可以实时采集用户所需要的模型里面或者用到的特征。之后会把实时特征平台跟实时计算平台做进行合并,用户最后只需要写SQL就可以实现所有的实时作业实现。
3.Jstorm全面启用Cgroup
当前由于部分历史原因导致现在很多作业跑在Jstorm上面,因此出现了资源分配不均衡的情况,之后会全面启用Cgroup。
4.在线模型训练
携程部分部门需要实时在线模型训练,通过用Spark训练了模型之后,然后使用Spark Streaming的模型,实时做一个拦截或者控制,应用在风控等场景。
—end—
- ?
不懂代码,如何做出实时刷新的数据大屏?
春天的
展开
首先恭喜你,当你看到这篇文章的时候,不管你是小白还是大咖,你都将直接获得一个高级技能:轻松上手可实时刷新的酷炫大屏。
制作可视化大屏,一般有这么几种方案:
写代码调用数据和图表,比如写JS+Echarts ;直接的数据可视化工具
前者对于大部分人来说门槛较高,而且尤其是大屏需求比较多,比方说要做10个的情况下,亲身试验写代码容易崩溃。如果涉及大量的动态可视化,涉及大数据量,没有底层技术,性能就会大打折扣。而且投到不同尺寸的屏幕,调试起来非常麻烦。
那么有没有一种简单的可视化大屏方案,可以快速的设计样式呈现效果、自适应不同大小的屏幕、而且还可以实时刷新数据?
有,选择后者,直接用数据可视化工具。
市面上能做到直接呈现在LED屏幕的大屏可视化工具并不多,多数需要代码调试,报表工具FineReport和FineBI工具可直接实现,相对来讲FineBI使用更简单,本文也是基于FineBI,来教大家做可实时刷新的数据大屏。
先来看看我们今天即将要教大家做的大屏效果(请接受一波酷炫可视化的冲击!)
不懂代码,如何做出实时刷新的数据大屏?
1、快速上手学习BI工具
FineBI是一个可视化的自助式BI工具,整个操作就是导数据/连数据库——处理数据(可视化ETL)选择图表——拖数据字段——可视化展现&美化,操作简单上手快。多数情况下,这个工具都是拿来做可视化报表,对接企业大数据平台,做企业数据运营分析用。
关于他的入门教程,小编之前曾发过一个视频《30分钟,教你零基础用BI搭建可视化大屏!》
2、构建数据模型
掌握了finebi的基础功能:怎么连接数据,怎么趋势,怎么做图表。接下来就到了正式做大屏步骤,先是构建数据模型。
大屏也是有主题的,本质是对一类业务的分析,然后综合展示,比如销售大屏。像这类业务分析一般要用到多张维度表和事实明细表的数据(例如下图中的分公司维度表和合同事实表)。常规操作是将不同业务系统的sql表拼接、宽表拼接,构成一个星型数据模型,需要你有专业的数据仓库技能。那这里化繁为简,可以直接用工具自带的敏捷数据模型去替代上述的工作,原理是自动构建雪花型模型,跨数据源关联。
搭建好上图的销售demo业务包的数据表和关联模型之后,下一步就可以进行正式的销售管理驾驶舱大屏搭建。
3、大屏布局设计
在给大家介绍具体制作过程之前先讲解一下通常管理驾驶舱的布局方式。管理驾驶舱往往展现的是一个企业全局的业务,一般分为主要指标和次要指标两个层次,主要指标反映核心业务,次要指标用于进一步阐述分析。所以在制作时给予不一样的侧重,这里推荐几种常见的版式。
上面几个版式不是金科定律,只是通常推荐的主次分布版式,能让信息一目了然。实际项目中,不一定使用主次分布,也可以使用平均分布,或者可以二者结合进行适当调整。比如下图所示,指标很多很多,存在多个层级的,就根据上面所说的基本原则进行一些微调,效果会很好。
4、实际分析制作过程
有了以上的布局设计,每一个模块就单独用一类图表分析一块内容,比如销售分布、签单分布、回款金额分析......整体呈现一个主题(在这里是销售业务)的分析。
那具体如何用工具操作呢?
首先,既然是销售管理驾驶舱,那么我们可以先从领导和高层最为关注的公司签单金额和回款金额入手。对于这样的汇总指标,选择仪表板进行数据展示再合适不过了。选择拖入合同事实表中的合同金额和合同回款表中的回款金额两个指标,样式这里选择圆环仪表盘,同时两个指标的单位都设置成亿,最大刻度输入当前合同金额,2.78亿。这样一来,2.78亿的合同回款,2.25亿的回款金额以及80.87%的总的回款率也就统计出来了,企业的签单金额和回款金额/回款率都一目了然。
其他部分也是一样的原理,篇幅原因不多介绍,核心是要知道展现哪些数据指标。
5、实时刷新功能
如何做出实时刷新的数据大屏,本篇还有一个重点内容就是大屏的实时刷新功能,也是大家问得比较多的。
所谓实时刷新,即你展示出来的酷炫大屏上面的数据将是动态刷新,能够实时反映数据库中的数据。我们的大屏通常连接着数据库,我们打开报表的时候,会读取数据库中的数据,但数据库中的数据可能是动态变化的,如果要读取变化的数据的话,不需要我们重新打开刷新报表,报表中的数据将动态自动刷新。
FineBI实时刷新的底层技术和性能:
实时刷新的实现所依靠的一个重要支撑,是FineBI自带的FineDirect直连引擎。FineDirect直连引擎给出了数据端到应用端的完整解决方案,支持连接企业已有的大数据计算平台,如Hadoop、Kylin、Greenplum、Vertica等,在充分利用平台计算性能的同时,也解决了TB至PB级超大数据量多维分析的难题。
FineDirect是FineBI推出的大数据直连引擎功能模块,用于更好地处理超大数据量的分析要求和数据源实时性的需求。通过FineDirect直连引擎可以直接对接现有的数据源,无论是传统的关系型数据库(Oracle,Sqlserver),还是日益成熟的Hadoop生态圈,Mpp架构的解决方案,都可以直接进行自助取数分析,实现更敏捷的、更及时的决策分析。
FineDirect引擎核心特点
①PB级别数据量多维分析
FineDirect直连引擎给出了数据端到应用端的完整解决方案,支持连接企业已有的大数据计算平台,如Hadoop、Kylin、Greenplum、Vertica等,在充分利用平台计算性能的同时,也解决了TB至PB级超大数据量多维分析的难题。
②实时大数据分析
FineDirect能够连接实时数据进行分析,及时返回分析结果。基于FineDirect的可视化引擎,可以将用户拖拽分析的操作,实时地转化为经过处理的查询语言,实现对企业数据库实时分析的效果。
③双引擎模式灵活搭配
FineBI已有FineIndex引擎(原cube)和新的FineDirect直连引擎可以搭配使用,来满足不同的应用场景。企业可以根据实际需求的不同准备两种类型的数据,通过FineIndex模式配置那些不经常更新、实时性要求不高的数据;通过FineDirect直连引擎配置大数据量且有实时分析需求的数据,双管齐下。
- ?
实时数据采集 将平台资源匹配
赵访枫
展开
匹配算法实现,通过将用户属性、用户学习状态属性、用户操作行为统计属性与资源属性进行匹配,数据容量大,遥测距离远,人机界面友好,可靠性高的优点,可广泛用于学校和社区小区域范围环境服务。
也适合于气象、海洋、环境、机场、港口、工及等领域使用.实时数据采集与上传该模块主要功能是通过无线方式获取环境观测站采集的实时数据并显示.数据采集的基本算法如下:首先判断是否到发送数据时间,若到了执行数据发送命令。
为了方便学习者查找资源以及系统推送资源,须对学习者和资源的属性进行标记,即创建属性元数据.本研究采用标签技术分别建立学习者和资源的属性元数据.学习者的属性标签主要来自于用户注册时的信息(用户属性)、学习者的学习状态属性以及学习者系统操作行为统计属性(对上传、下载、浏览、评论等操作的统计)。
资源的属性标签主要来自于资源提供者根据学科和关键字分类的标签以及资源使用者对资源的评价、标记等.根据研究性学习的主题、用户学习进度,系统可以将用户可能感兴趣,其他用户感兴趣或评分高、浏览次数多,相同领域的最优资源或相似领域的最佳资源推送给用户,系统根据学习者的学习主题,将平台中评价最高的相关资源、教师指定的资源和拓展的相关资源进行匹配,
- ?
回顾·大数据平台从0到1之后
远锋
展开
本文根据链家赵国贤老师在DataFun Talk数据架构系列活动“海量数据下数据引擎的选择及应用”中所分享的《大数据平台架构从0到1之后》编辑整理而成,在未改变原意的基础上稍做修改。
大数据平台构建方法大同小异,但是平台构建以后也面临很多挑战,在面临这些挑战我们如何去克服、修复它,让平台更好满足用户需求,这就是本次主题的重点。下面是本次分享的内容章节,首先讲一下架构1.0与2.0,两者分别是怎么样的,从1.0到2.0遇到了哪些问题;第二部分讲一下数据平台,都有哪些数据平台,这些数据平台都解决什么问题;第三个介绍下当前比较重要的项目“olap引擎的选型与效果”以及遇到的一些问题;第四个简单讲一下在透明压缩方面的研究。
架构1.0阶段,底层是Hadoop,用来存储数据和分析数据。需要把log数据和事务数据传输到Hadoop平台上,我们使用的是kafka和sqoop进行数据传输。然后在Hadoop平台基础上,通过一个开源的Hive和oozie做一个调度,开发者写Hql来完成业务需求,然后将数据mysql集群或redis集群,上层承接的是一个报表系统。这个需求基本跑了一年,也解决了一些问题。但存在的问题有:(1)架构简单,不易解耦,结合太紧密出现问题需要从底层一直查到上面;(2)平台架构是需求驱动,面临一个需求后需要两周时间来解决问题,有时开发出来运营已经不需要;(3)将大数据工程师做成一个取数工程师,大量时间在获取怎样数据;(4)故障频发,比如Hql跑失败了或者网络延迟没成功,oozie是通过xml配置发布任务,我们解决需要从数据仓库最底层跑到数据仓库最高层,还要重刷msl,花费时间。
面对这些问题我们做了一次架构调整,数据平台分为三层,第一层就是集群层(Cluster),主要是一些开源产品,Hadoop实现分布式存储,资源调度Yarn,计算引擎MapReduce、spark、Presto等,在这些基础上构建数据仓库Hive。还有一些分布式实时数据库HBase还有oozie、sqoop等,这些作用就是做数据存储、计算和调度,另外还有一个数据安全。第二层就是工具链,这一层是一个自研发调度平台,架构1.0用的oozie。基本满足需求有调度分发,监控报警,还有智能调度、依赖触发,后续会详细介绍。出问题后会有一个依赖关系可视化,数据出问题可以很快定位与修复。然后就是Meta(元数据管理平台),数据仓库目前有3万多张表,通过元数据管理平台实现数据仓库数据可视化。还有一个AdHoc,将数据仓库中的表暴露出去,通过平台需求方就可以自主查找自己需要的数据,我只需要优化查询引擎、记录维护、权限控制、限速和分流。最上层将整个大数据的数据抽象为API,分为三个,面向大数据内部的API,面向公司业务API,通用API。大数据内部API可以满足数据平台一些需求,如可视化平台、数据管理平台等,里面有专有API来管理这些API。面向公司业务API,我们是为业务服务的,通过我们的技术让业务产生更多产出,将用户需要的数据API化,通过API获取数据就行。通用API,数据仓库内部的报表都产生一些API,业务需求方根据自己的需求自动组装就OK了。架构2.0基本解决了我们架构1.0解决的问题。
第二部分就简单介绍下平台,第一个是存储层-集群层,解决运维工作,我们基于开源做了一个presto。实习人员经过一两周能适应这个工作,释放了运维的压力,数据量目前有18PB,每天的任务有9万+,平均3-4任务/分钟;第二个就是元数据管理平台,这种表抽象为各个层,分析数据、基础细节数据等抽象,提供一个类似百度的搜索框,通过搜索获得所需数据,这样业务人员能够非常方便的使用我们的数据。它能实现数据地图(数据长怎样,关联关系是怎么样都可以显示出来),数据仓库可视化,管理运维数据,数据资产非常好的管理和运维,将数据开发的工作便捷化、简易化。
第三个数据平台调度系统,数据仓库中的各个层需要流转,数据出现问题后如何去恢复数据。数据调度系统主要的工作有:(1)数据流转调度,可以非常简易的配置出数据的流转调度。(2)依赖触发,充分利用资源,能够让调度任务非常紧凑,能够尽可能快的产出我们的数据。(3)对接多个数据源,需要将多种多样的数据源集成到数据仓库中,如何将sql server数据、Oracle数据等数据导入到数据仓库中,系统能够对接多种数据源,因此我们财务人员、运营人员、业务人员都可以自主将数据接入到数据仓库,然后分析和调度。(4)依赖关系可视化。比如我们有100个任务是关联的,最底层std层有50个任务,中间层有20个任务,如果中间ODS层出问题了,会影响上层依赖层任务,通过可视化就能很方便定位。
除了前面三个平台,还需要一个平台来展示我们的数据,才能向我们的用户显示数据的价值。我们的指标平台支持上卷下钻、多维分析、自助配置报表,统一公司的各个指标。说一下统一公司的各个指标,比如链家场景,比如说一个业绩(一周卖出十套房子,需要提佣),16年我们发现有多个口径,因此通过指标系统将指标统一化,指标都从这里出,可以去做自己的可视化。还有各种财务人员、区长或店长也可以自主从指标平台上配置自己的数据,做自己的desktop,指标系统的后端使用后续讲Kylin的一个多维分析引擎支撑的。
指标平台架构,一个应用的可视化平台肯定需要底层能力的支撑,这次主题也是数据引擎,链家使用的是一个叫kylin的开源数据引擎,可以把数据仓库中的数据通过集群调度写入到HBase中做一个预计算。这样就可以支持指标系统千亿级数据亚秒级的查询,不支持明细查询因为做过预计算。还引入了百度开源的palo,经过优化,通过这样一个架构就满足上层的地动仪、指标平台和权限系统。运营、市场、老板都在用这个指标平台,能够实现多维分析、sql查询接口、超大规模数据集、释放数据的能力以及数据可视化。
我们是需求驱动,每天都会遇到很多需求,数据开发人员就是取出需要的数据。利用adhoc平台将数据从数据仓库中取出,基于这个我们做了一个智能搜索引擎,架构在adhoc上的搜索引擎有很多,比如presto、hive、spark等。用户也不知道该选择那种引擎,他的需求就是尽可能取出自己所需的数据,因此开发智能选择引擎、权限控制,并且能够支撑各种接口、自助查询,这样就基本解决了数据开发的工作。我们自研发了一个queryengine,在底层有presto、sparksql、hive等,queryengine特点就是能够发挥各自引擎的特性,如presto查询快,但是sql支撑能力不强,sparksql同样,在某些特殊sql查询不如hive快,hive就是稳但是慢。queryengine就是智能选择各种引擎,用户把sql提交过来,queryengine判断哪个引擎适合你。如何做的简单介绍下,对sql进行解析成使用的函数、使用的表、需要返回的字段结构,根据各个引擎的能力判断哪个合适。目前还在开发功能就是计费,因为资源是有限的。queryengine支持mysql协议,因为有些用户需要BI能力,需要对返回的数据进行聚合,我们不能开各种各样的BI能力,我们只需满足mysql协议将数据暴露出去,用户只需用其他BI就能使用。
通过架构1.0到架构2.0衍生出很多平台,大架构已经有了,但是遇到的一些问题如何解决。这里分享两个案例,一个是olap引擎的选型与效果,第二个就是为什么要做透明压缩,是如何做的。Rolap引擎基本是基于关系型数据库,基于关系模型实时进行聚合运算,主要通过传统数据库或spqrk sql和presto,spqrk sql和presto是根据数据实时计算;Molap是基于一个预定义模型,预先进行聚合计算,存储汇总结果。先计算好一个立方体,基于立方体做上传下钻,实现由Kylin/Druid,Druid主要是实时接入(Kylin没有),实时将kafka数据用Spark sql做一次计算然后将数据上传上去,可以支持秒级查询;还有一个比较流行的是叫olap,混合多引擎,不同场景路由到不同引擎。
Rolap查询时首先将数据扫描出来,然后进行聚合,通过聚合结果将多个节点数据整合到一个节点上然后返回。优势是支持任何sql查询,因为数据是硬算,使用明细数据,没有数据冗余,一致性非常好,缺点是大数据量或复杂数据量返回慢,因为你是基于明细数据,一条一条数据计算无论如何优化还是会出现瓶颈,并发性很差。
Molap中间会有一个中心立方体cube,在数据仓库通过预计算将数据存储到cube中,通过预聚合存储支持少量计算汇总,为什么少量计算,因为数据都已经预计算好了。优点就是支持超大数据集,快速返回并发高,缺点是不支持明细,需要预先定义维度和指标,适用场景就是能预知查询模式,并发有要求的场景,固化场景可以使用molap。
对于技术选型,当时面临的需求,基本上开源组件有很多,为什么选择kylin,因为支持较高的并发,面对百亿级数据能够支持亚秒级查询,以离线为主,具有一定的灵活性,最好有sql接口,而这些需求刚好kylin能满足。Apache Kylin是一个开源的分布式分析引擎,提供Hadoop之上的SQL查询接口及多维分析能力,以支持超大规模数据,最初由e Bay Inc. 开发并贡献至开源社区。它能在亚秒内查询巨大的Hive表。其解决方案就是预先定义维度和指标,预计算cube,存储到hbase中,查询时解析sql路由到hbase中获取结果。
现在讲一下链家olap架构,HBase集群,数据仓库计算和预处理在这块,还有一个为了满足kylin需求而做的HBase集群。Kylin需要做预计算,因此有个build集群,将数据写入到基于kylin的Hadoop集群中,然后利用nginx做一个负载均衡,还有一个query集群,然后就是面向线上的一个查询,还有一个kylin中间件,解决查询、cube任务执行、数据管理、统计。指标平台大部分是查询kylin,但是kylin不能满足明细查询,这个就通过queryengine智能匹配,通过spark集群或presto集群,还有alluxio做压缩,然后将明细查询结果返回指标平台,最终返回其他业务的产品。在横向还做了一个权限管理、监控预警、元数据管理、调度系统,来实现整体平台支撑。
接下来讲一下链家kylin能力拓展,基本大同小异,遇到的问题主要有:分布式构建,cube增长很快,build集群无法承载,因此做了分布式优化能够满足500cube在规定时间跑完;优化构建时字典下载策略,kylin构建时需要将所有元数据字典全部下载下来,因此从Hadoop将元数据字典下载都得好几分钟,每次build都去下载元数据字典会很耗时,优化后只需要下载一次就可以;优化全局字典锁,build时需要锁住整个build集群,完成后锁才释放,源码发现并不需要全局锁只需要锁住所需要的字段就可以,优化将锁设置到字段级别上;Kylin 的query查询机器使用G1垃圾回收器。我们自研发了一个中间件基本可以容纳一个无限容量的队列,针对特定cube的预先调度,以及权限的管控、实现任务的并发控制。架构有外面的调度系统,有一个kylin中间件,所有的查询和build都经过kylin中间件。还做了一个任务队列、统计、优先级调度、监控报警、cube平分、以及可视化配置和展示。
架构从0到1.0遇到了另一个问题-集群,存储链家所有数据,数据量大、数据增长快(0-1PB两年时间,1PB-16PB不到一年时间,面临成本问题)、冷数据预期,针对这些问题提出透明压缩项目。就是分层存储(Hadoop特性),根据不同数据分不同级别存储,比如把一部分数据存储在ssd,把另一部分数据存储到磁盘之上。Hot策略将数据全部存储到磁盘之上,warm策略就是一部分数据存储在磁盘上,一部分存储archive(比较廉价,转数小)。第二个就是ZFS文件系统,它具有存储池、 自我修复功能、压缩与可变块大小、 写时拷贝/校验和/快照、 ARC(自适应内存缓存)与L2ARC(SSD做二级缓存)。
透明压缩设计实现思路是:(1)界定要做数据冷处理隔离的主要内容。需要将一部分数据存储到ZFS文件系统做一个透明压缩来满足减少成本的需求,这样需要把冷数据界定出来;(2)生成特定的通过获取特定的冷数据列表,并标记其冷数据率;然后,定期从冷数据表中取出为完成冷数据迁移的行,进行移动。通过HDFS目录把界定出来的冷数据移动到ZFS压缩之上,把不需要的移除到Ext4上。这样一部分数据存储在ZFS上,一部分存储在EXT4上。
透明压缩优化工作有:第一个Hadoop冷热数据分离优化。涉及有异构存储策略选择、HDFS冷热数据移动优化;第二个就是ZFS文件系统优化。ZFS支持很多压缩算法,经过测试发现Gz压缩效率最好,下图是各种算法效率对比。随着压缩数据越来越大,CPU占用越来越高。海量数据集群不光是存储还有计算。Datanode对压缩数据的加载时间,直接关系到访问此部分数据时的效率,从表可知,ZFS的gz压缩在datanode加载数据上对LZ4有部分优势。较为接近EXT4。综合考虑压缩率,读取,写入速度,datanode加载速度等,选定gz作为ZFS文件系统的压缩算法。
透明压缩前数据增长是非常快的,接近30%的增长速率,逻辑数据有3PB,3备份后总空间:9.3PB实际总空间:7PB,就目前简单预估节省成本有300万。压缩后虽然实际数据再增长,但真实数据是缓慢下降的。
透明压缩未来展望,透明压缩是对cpu是有损耗的,我们希望将透明压缩计算提取出来,通过QAT卡进行压缩,希望...
- ?
5个好用的可视化数据平台,让你的数据分析更高效率、高逼格
纪寇
展开
文|小天文章源自:起点学院学习联盟
在小白们眼里,大神们的数据分析报表基本上是这样的……
要么就是像这样的……
而大部分人,差不多是这样的……
啊多么痛的领悟……
怎样才能又快又好地做出一份高颜值的数据报表呢?带着立志要把这样的图表从癞蛤蟆脱胎成白天鹅的坚定和悲壮,这里搜集了5个笔者之前用过,用户评价不错,用起来还顺手的可视化数据平台。
话不多说,直接上正文。如果你也有推荐的平台,欢迎留言分享~
-1-Echarts
没想到这个第一次用就惊艳到我的产品竟然是国产,而且还来自百度,简直堪称良心。先上几张用Echarts制作的效果图。
貌似很多小伙伴喜欢用Echarts制作地图类的可视化效果……毕竟酷炫……
除了这些惊艳的地图,Echarts同样可以运用于散点图、折线图、柱状图等这些常用的图表的制作。
如果你需要展示实时变化的数据,相信Echarts里的动态接口会对你十分有帮助。
Echarts的优点在于,文件体积比较小,打包的方式灵活,可以自由选择你需要的图表和组件。而且图表在移动端有良好的自适应效果,还有专为移动端打造的交互体验。
-2-Highcharts
这个也是很多小伙伴在使用的一个平台。完全不用担心找不到参考的样图,因为已经有很多中国区的用户在上面更新并维护着很多实例,你往往能从这些丰富的例子找到类似的表达样图。
它的图表类型自然也是很丰富啦,线图、柱形图、饼图、散点图、仪表图、雷达图、热力图、混合图等类型的图表都可以制作,也可以制作实时更新的曲线图。
Highcharts对非商用免费,对于个人网站,学校网站和非盈利机构,可以不经过授权直接使用Highcharts系列软件。Highcharts还有一个好处在于,它完全基于HTML5技术,不需要安装任何插件,也不需要配置PHP、Java等运行环境,只需要两个JS文件即可使用。
-3-帆软报表(FineReport)
FineReport的可视化效果虽然没有上面两种那么酷炫,因为定位是报表软件。但是赢在操作相当简易,不会上面那些复杂的代码也没关系。它采用类似于Excel的编辑器,只需要点选拖拽等操作,拖动数据列绑定至对应单元格,简单设置就可以在web端查看数据展示。
目前,它有普通报表、聚合报表和决策报表三类报表设计模式,基本可以满足企业各类日常数据分析的情景需求。
数据的可视化与交互效果也很不错。最牛逼的可以做高大上的动态报表
还有一个比较强大的地方,就是它的数据填报。区别于传统意义上只能做数据展示的报表,FineReport允许用户对数据库的增删改。而且,它填报报表的流程非常简单,只要四步:报表设计、控件添加、设置填报属性和填报录入,这样,填报工作就能轻松搞定啦~
-4-数说立方
数说立方是大数据应用与服务提供商“数说故事”旗下一款面向数据分析师的在线商业智能产品。在数据的可视化呈现方面,操作比较简便,即使是非数据分析的专业人员,也能轻松实现。
同时,它的实时数据可视化引擎也能让使用者可以第一时间获得数据的可视化反馈,直观地了解到数据的变化情况。
-5-PowerBI
PowerBI是微软发布的一款可视化BI工具,类似Excel升级版的大表哥。一改以往excel需要数据透视表,写大量函数的复杂特点,这款工具拖拖拽拽操作起来十分简单。
- ?
2018开源数据库论坛(ODF)议:如何设计实时数据平台-下
幻影
展开
2018开源数据库论坛(ODF)议题:如何设计实时数据平台-下篇(技术篇)
导读:实时数据平台(RTDP,Real-time Data Platform)是一个重要且常见的大数据基础设施平台。在上篇(设计篇)中,我们从现代数仓架构角度和典型数据处理角度介绍了RTDP,并探讨了RTDP的整体设计架构。本文作为下篇(技术篇),则是从技术角度入手,介绍RTDP的技术选型和相关组件,探讨适用不同应用场景的相关模式。RTDP的敏捷之路就此展开~
文章链接:
https://mp.weixin.qq/s?__biz=MzU0MDExOTUyMg==&mid=2247484212&idx=1&sn=21d979b484942e193bac15032abbdb01&chksm=fb3f5bb9cc48d2afb6e53f6a2850ac70ef57a5d90f542745611556f86a537f2207c420bb22df#rd
2018 ODF 开源数据库论坛暨MariaDB中国用户者大会演讲嘉宾预告
卢山巍9月8日,宜信大数据技术专家卢山巍将在“2018 ODF 开源数据库论坛暨MariaDB中国用户者大会”的大数据专场作主题演讲,分享《敏捷大数据实践与开源赋能》,主要介绍敏捷大数据思想、方法学、平台,以及支撑的敏捷大数据实践应用。主要围绕实时数据平台这个话题展开,从痛点到架构再到实践应用,给出一个端到端的实时数据平台解决方案。同时也会介绍架构中我们研发的开源平台。
欢迎来现场与大咖交流。
ODF官网: http://odf
大会百格活动:
https://bagevent/event/1560938
实时数据库平台
-
1、只需3秒快速实现求和
-
2、如何快速填充序号
-
3、如何自动填充序号(公式法)
-
4、数据条的神奇应用
-
5、多文本快速合并
-
6、查找与替换的不同玩法
-
7、快速定位到指定区域
-
8、数据排序、工资条制作
-
9、快速筛选(模糊、精确筛选)
-
10、快速插入空行
-
11、快速删除空行
-
12.快速跳转到天涯海角
-
13、.同时查看两个Excel文件
-
14、用条件格式扮靓报表
-
15、一键插入Excel图表
-
16、批量处理行高、列宽
-
17、利用拆分功能查看数据
-
18、批量录入相同内容
-
19、工作表快速跳转
-
20、批量录入表格模板(精品课程)
-
21、Excel函数与公式的应用、公式循环引用的查找
-
22、IF函数单条件判断同比增长
-
23、用sum函数 格式相同,连续多表数据汇总
-
24、excel快捷键
-
25、VLOOKUP函数——根据销售员匹配销售额
-
26、统计各部门销售总额
-
27、统计指定条件个数
-
28、怎样输入当前日期和时间、星期数
-
29、销售业绩排名
-
30、Sumproduct函数-万能函数(销售额汇总求和)
-
31、根据销售员,地区,商品名称汇总
-
32、批量替换PPT字体
-
33、给销售额数据批量添加万元单位
-
34、一秒快速核对两列数据
-
35、快速定位到指定单元格或区域
-
36、快速制作双行标题工资条
-
37、给你的表格做个瘦身
-
38、快速打开常用的Excel文件
-
39、快速打开多个Excel文件
-
40、利用创建组—快速隐藏/展开多列数据
-
41、快速制作下拉菜单
-
42、复制粘贴表格,如何保留数据源列宽格式一致?
-
43、两列数据位置互换
-
44、1秒钟扮靓报表——如何实现表格隔行换色
-
45、快速删除重复记录——保留唯一值
-
46、快速向下填充、向右填充,文本或公式
-
47、给Excel文件添加密码
-
48、插入带图片的批注
-
49、输入公式后不计算?
-
50、如何设置单元格缩进
-
51、快速解决Excel表格总显示货币格式
-
52、批量添加万元单位
-
53、你会四舍五入么?
-
54、用RAND函数机选彩票
-
55、冻结首行你会么?
-
56、超链接的高级应用
-
57、IFERROR函数-屏蔽错误值
-
58、批量填充颜色
-
59、录入数据
-
60、快速输入工号
-
61、快速行列转置
-
62、自定义缩放界面
-
63、多个单元格同时输入
-
64、如何计算立方米?
-
65、快速制作双行标题工资条
-
66、输入带方框的√和×
-
67、快速将姓名对齐
-
68、快速输入性别
-
69、按单位职务排序
-
70、自动计算合同到期日期
-
71、计算时间间隔
-
72、日期和时间的拆分
-
73、快速处理不规范的日期格式
-
74、快速填充合并单元格
-
75、效率加倍的快捷键
-
76、快速复制表格和对象
-
77、快速创建工作表副本
-
78、快速复制序列号
-
79、快速显示公式
-
80、多个单元格同时输入
-
81、快速调整显示比例
-
82、快速自动填充
-
83、快速填充(Ctrl+E)
-
84、Ctrl与数字键结合
-
85、快速将多列数据整理为1列
-
86、快速将1列数据拆分为多列
-
87、快速定位公式
-
88、快速录入数据
-
89、快速累计求和
-
90、身份证号码显示为0怎么办?
-
91、快速制作斜线表头
-
92、文本竖向显示
-
93、神奇的监视窗口
-
94、不一样的格式刷
-
95、快速美化图表
-
96、快速生成当前日期
-
97、快速找出循环引用
-
98、快速提取信息
-
99、二维表快速转换为一维表
-
100、快速多表合并