[第十一章]流计算

概述

静态数据和流数据

  • 静态数据
    • 不会随时间变化的数据
  • 流数据
    • 时间分布和数量无上限的一系列动态数据集合体
    • 数据记录是最小组成单元
    • 特征
      • 数据快速持续到达,潜在大小也许是无穷无尽的 数据来源众多,格式复杂 数据量大,但是不十分关注存储,一旦经过处理,要么被丢弃,要么被归档存储 注重数据的整体价值,不过分关注个别数据 数据顺序颠倒,或者不完整,系统无法控制将要处理的新到达的数据元素的顺序

批量计算和实时计算

  • 静态数据批量计算
  • 流数据实时计算
  • 需求
    • 高性能
    • 海量式
    • 实时性
    • 分布式
    • 易用性
    • 可靠性

流计算概念

实时获取来自不同数据源的海量数据,经过实时分析处理,获得有价值的信息
数据价值随时间流逝而降低
 

流计算与Hadoop

Hadoop擅长批处理,不适合流计算

流计算框架

  1. 商业级
    1. IBM InfoSphere Streams
    2. IBM StreamBase
  1. 开源
    1. Twitter Storm
    2. Yahoo! S4
  1. 公司自己框架
    1. DStream
    2. Super Mario

流计算的处理流程

概述

  • 传统数据计算流程存在问题
    • 数据是旧的
    • 用户主动发起查询

数据实时采集

  • 实时性
  • 低延迟
  • 稳定可靠

基本框架

Agent:主动采集数据,并把数据推送到Collector部分 Collector:接收多个Agent的数据,并实现有序、可靠、高性能的转发 Store:存储Collector转发过来的数据(对于流计算不存储数据
notion image

数据实时计算

实时查询服务

流计算的应用

开源流计算框架Storm

特点

整合性:Storm可方便地与队列系统和数据库系统进行整合 简易的API:Storm的API在使用上即简单又方便 可扩展性:Storm的并行特性使其可以运行在分布式集群中 容错性:Storm可自动进行故障节点的重启、任务的重新分配 可靠的消息处理:Storm保证每个消息都能完整处理 支持各种编程语言:Storm支持使用各种编程语言来定义任务 快速部署:Storm可以快速进行部署和使用 免费、开源:Storm是一款开源框架,可以免费使用

设计思想

  • Streams
    • Strom将流数据Stream描述成一个无限的Tuple序列,这些Tuple序列会以分布式的方式并行地创建和处理
    • notion image
  • Tuple
    • 每个tuple是一堆值,每个值有一个名字,并且每个值可以是任何类型
  • spout
    • stream的来源
    • 通常Spout会从外部数据源(队列、数据库等)读取数据,然后封装成Tuple形式,发送到Stream中。nextTuple函数,Storm框架会不停的调用该函数
  • Bolt
    • Streams的状态转换过程。可以执行过滤、函数操作、Join、操作数据库等任何操作
    • 其接口中有一个execute(Tuple input)方法,在接收到消息之后会调用此函数
  • Topology
    • Spouts和Bolts组成的网络抽象成Topology
    • 组件之间的连接则表示数据流动的方向
    • 并行运行的
  • Stream Groupings
    • 用于告知Topology如何在两个组件间进行Tuple的传送
    • 决定传送方式和时间
    • 类型
      • ShuffleGrouping:随机分组,随机分发Stream中的Tuple,保证每个Bolt的Task接收Tuple数量大致一致
      • FieldsGrouping:按照字段分组,保证相同字段的Tuple分配到同一个Task
      • AllGrouping:广播发送,每一个Task都会收到所有的Tuple
      • GlobalGrouping:全局分组,所有的Tuple都发送到同一个Task中
      • NonGrouping:不分组,和ShuffleGrouping类似,当前Task的执行会和它的被订阅者在同一个线程中执行
      • DirectGrouping:直接分组,直接指定由某个Task来执行Tuple的处理

框架设计

  • MapReduce作业最终会完成计算并结束运行,而Topology将持续处理消息(直到人为终止)
notion image

Spark Streaming

  • 数据流以时间片(秒级)为单位进行拆分,然后经Spark引擎以类似批处理的方式处理每个时间片数据
  • DStream
    • 离散化数据
    • 按照时间片(如1秒)分成一段一段的DStream,最终转变为对相应的RDD的操作

对比⭐

  • Spark Streaming无法实现毫秒级的流计算,而Storm可以实现毫秒级响应
  • Spark Streaming构建在Spark上,相比于Storm,RDD数据集更容易做高效的容错处理
  • Spark Streaming采用的小批量处理的方式使得它可以同时兼容批量和实时数据处理的逻辑和算法

Samza

概念

  • 作业
    • 一组输入流进行处理转化成输出流的程序
  • 分区
    • 分区是一个有序消息序列,消息是Samza的流数据的单位
    • 每个流都被分割成一个或多个分区
  • 任务
    • 作业会被进一步分割成多个任务(Task)来执行,其中,每个任务负责处理作业中的一个分区
    • YARN调度器负责把任务分发给各个机器
  • 数据流图
    • 一个数据流图是由多个作业构成的,其中,图中的每个节点表示包含数据的流,每条边表示数据传输
    • Job、Stream

架构

  • 主要包括流数据层(Kafka)、执行层(YARN)、处理层(Samza API)
notion image

Strom、Spark Streaming、Samza应用场景

  • 从编程的灵活性来讲,Storm是比较理想的选择
  • 需要在一个集群中把流计算和图计算、机器学习、SQL查询分析等进行结合时,可以选择Spark Streaming
  • 当有大量的状态需要处理时,比如每个分区都有数十亿个元组,则可以选择Samza
 
Prev
[第十章]Spark
Next
[第十二章]Flink
Loading...
Article List
一个NotionNext搭建的博客
数据库系统概论
大数据原理与应用
javaWeb应用开发基础教程
python
毕业设计
大数据技术综合应用
实训-航空数据系统
java面向对象程序设计
数据结构
算法分析与设计
SPARK
Python爬虫大数据采集与挖掘
云计算
概率论与数理统计
数字逻辑
计算机网络
计算机组成原理
linux
操作系统
人工智能导论
数据仓库与数据挖掘
数据可视化
大数据安全与隐私保护
c语言
C++