Apache Flink

Apache Flink是由Apache软件基金会开发的开源流处理框架,其核心是用Java和Scala编写的分布式流数据流引擎。Flink以数据并行和管道方式执行任意流数据程序,Flink的流水线运行时系统可以执行批处理和流处理程序。此外,Flink的运行时本身也支持迭代算法的执行。

Flink提供高吞吐量、低延迟的流数据引擎以及对事件-时间处理和状态管理的支持。Flink应用程序在发生机器故障时具有容错能力,并且支持exactly-once语义。程序可以用Java、Scala、Python和SQL等语言编写,并自动编译和优化到在集群或云环境中运行的数据流程序。

Flink并不提供自己的数据存储系统,但为Amazon Kinesis、Apache Kafka、Alluxio、HDFS、Apache Cassandra和Elasticsearch等系统提供了数据源和接收器。

开发
Apache Flink是由Apache软件基金会内的Apache Flink社区基于Apache许可证2.0开发的,该项目已有超过100位代码提交者和超过[https://github.com/apache/flink 460贡献者。]

[http://www.data-artisans.com/ data Artisans]是由Apache Flink的创始人创建的公司。目前,该公司已聘用了12个Apache Flink的代码提交者。

概述
Apache Flink的数据流编程模型在有限和无限数据集上提供单次事件(event-at-a-time)处理。在基础层面,Flink程序由流和转换组成。 “从概念上讲,流是一种(可能永无止境的)数据流记录,转换是一种将一个或多个流作为输入并因此产生一个或多个输出流的操作”。

Apache Flink包括两个核心API:用于有界或无界数据流的数据流API和用于有界数据集的数据集API。Flink还提供了一个表API,它是一种类似SQL的表达式语言,用于关系流和批处理,可以很容易地嵌入到Flink的数据流和数据集API中。Flink支持的最高级语言是SQL,它在语义上类似于表API,并将程序表示为SQL查询表达式。

编程模型和分布式运行时
Flink程序在执行后被映射到流数据流。

状态:检查点、保存点和容错
Apache Flink具有一种基于分布式检查点的轻量级容错机制。用户可以生成保存点,停止正在运行的Flink程序,然后从流中的相同应用程序状态和位置恢复程序。 保存点可以在不丢失应用程序状态的情况下对Flink程序或Flink群集进行更新。从Flink 1.2开始,保存点还允许以不同的并行性重新启动应用程序,这使得用户可以适应不断变化的工作负载。

数据流API
Flink的[https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/datastream_api.html 数据流API]支持有界或无界数据流上的转换(如过滤器、聚合和窗口函数),包含了20多种不同类型的转换,可以在Java和Scala中使用。

有状态流处理程序的一个简单Scala示例是从连续输入流发出字数并在5秒窗口中对数据进行分组的应用:
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.windowing.time.Time

case class WordCount(word: String, count: Int)

object WindowWordCount {
def main(args: Array[String]) {

val env = StreamExecutionEnvironment.getExecutionEnvironment
val text = env.socketTextStream("localhost", 9999)

val counts = text.flatMap { _.toLowerCase.split("\\W+") filter { _.nonEmpty } }
.map { WordCount(_, 1) }
.keyBy("word")
.timeWindow(Time.seconds(5))
.sum("count")

counts.print

env.execute("Window Stream WordCount")
}
}

Apache Beam - Flink Runner
Apache Beam“提供了一种高级统一编程模型,允许(开发人员)实现可在在任何执行引擎上运行批处理和流数据处理作业”。Apache Flink-on-Beam运行器是功能最丰富的、由Beam社区维护的能力矩阵。

data Artisans与Apache Flink社区一起,与Beam社区密切合作,开发了一个强大的Flink runner。

数据集API
Flink的[https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/batch/index.html 数据集API]支持对有界数据集进行转换(如过滤、映射、连接和分组),包含了20多种不同类型的转换。 该API可用于Java、Scala和实验性的Python API。Flink的数据集API在概念上与数据流API类似。

表API和SQL
Flink的[https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/table_api.html 表API]是一种类似SQL的表达式语言,用于关系流和批处理,可以嵌入Flink的Java和Scala数据集和数据流API中。表API和SQL接口在关系表抽象上运行,可以从外部数据源或现有数据流和数据集创建表。表API支持关系运算符,如表上的选择、聚合和连接等。

也可以使用常规SQL查询表。表API提供了和SQL相同的功能,可以在同一程序中混合使用。将表转换回数据集或数据流时,由关系运算符和SQL查询定义的逻辑计划将使用Apache Calcite进行优化,并转换为数据集或数据流程序。

Flink Forward
[http://flink-forward.org/ Flink Forward]是一个关于Apache Flink的年度会议。第一届Flink Forward于2015年在柏林举行。为期两天的会议有来自16个国家的250多名与会者。 会议分为两个部分,Flink开发人员提供30多个技术演示,另外还有一个Flink培训实践。

2016年,350名与会者参加了会议,40多位发言人在3个平行轨道上进行了技术讲座。第三天,与会者被邀请参加实践培训课程。

2017年,该活动也将扩展到旧金山。 会议致力于Flink如何在企业中使用、Flink系统内部、与Flink的生态系统集成以及平台的未来进行技术会谈。它包含主题演讲Flink用户在工业和学术界的讲座以及关于Apache Flink的实践培训课程。

来自以下组织的发言人在Flink Forward会议上发表了演讲:阿里巴巴集团、、、第一资本、Cloudera、data Artisans、EMC、爱立信、Hortonworks、华为、IBM、Google、MapR、MongoDB、Netflix、、,Red Hat、ResearchGate、Uber和Zalando。

历史
2010年,研究项目“Stratosphere:云上的信息管理”(由德国研究基金会(DFG)资助)由柏林工业大学、柏林洪堡大学和哈索·普拉特纳研究院合作启动。Flink从Stratosphere的分布式执行引擎的一个分支开始,于2014年3月成为Apache孵化器项目。2014年12月,Flink成为Apache顶级项目。

发布日期

** 04/2015: [https://flink.apache.org/news/2015/04/13/release-0.9.0-milestone1.html Apache Flink0.9-里程碑-1]

Apache孵化器发布日期

Pre-Apache Stratosphere 发布日期

  • 01/2014: Stratosphere 0.4(0.3版本被跳过)
  • 08/2012: Stratosphere 0.2
  • 05/2011: Stratosphere 0.1(08/2011:0.1.1)

参见
*

  • 其他类似的数据处理引擎,如Storm和Spark。
  • Apache Beam,一种共享编程模型,Flink是其创始后端。

参考文献
外部链接
*

评论 (0)

  • 还没有评论,来抢沙发吧。