Apache Flink 是一个开源流处理框架,具有强大的流处理和批处理能力。
了解更多 Flink 信息,请访问 https://flink.apache.org/。
一个以流处理优先的运行时,同时支持批处理和流式数据程序
在 Java 和 Scala 中提供优雅流畅的 API
一个同时支持超高吞吐量和低事件延迟的运行时
在 DataStream API 中支持事件时间和乱序处理,基于 Dataflow Model
在不同时间语义(事件时间、处理时间)下提供灵活的窗口操作(时间、计数、会话、自定义触发器)
具有 精确一次 处理保证的容错能力
流式程序中自然的背压
用于图处理(批)、机器学习(批)和复杂事件处理(流)的库
在 DataSet(批)API 中内置对迭代程序(BSP)的支持
自定义内存管理,支持在内存(in-memory)和内存外(out-of-core)数据处理算法之间高效、稳健地切换
为 Apache Hadoop MapReduce 提供兼容层
与 YARN、HDFS、HBase 以及 Apache Hadoop 生态系统的其他组件集成
case class WordWithCount(word: String, count: Long)
val text = env.socketTextStream(host, port, '\n')
val windowCounts = text.flatMap { w => w.split("\\s") }
.map { w => WordWithCount(w, 1) }
.keyBy("word")
.timeWindow(Time.seconds(5))
.sum("count")
windowCounts.print()
case class WordWithCount(word: String, count: Long)
val text = env.readTextFile(path)
val counts = text.flatMap { w => w.split("\\s") }
.map { w => WordWithCount(w, 1) }
.groupBy("word")
.sum("count")
counts.writeAsCsv(outputPath)
构建 Flink 的先决条件:
git clone https://github.com/apache/flink.git
cd flink
mvn clean package -DskipTests # this will take up to 10 minutes
Flink 现在已安装在 build-target 中。
注意:Maven 3.3.x 可以构建 Flink,但无法正确遮蔽某些依赖项。Maven 3.1.1 可以正确创建这些库。 如果使用 Java 8 构建单元测试,请使用 Java 8u51 或更高版本,以避免使用 PowerMock 运行器的单元测试失败。
Flink 的提交者使用 IntelliJ IDEA 来开发 Flink 代码库。 我们推荐使用 IntelliJ IDEA 来开发涉及 Scala 代码的项目。
IDE 的最低要求是:
IntelliJ IDE 开箱即用地支持 Maven,并提供了用于 Scala 开发的插件。
有关详细信息,请参阅我们的 设置 IntelliJ 指南。
注意: 根据我们的经验,这种设置不适用于 Flink, 原因是 Scala IDE 3.0.3 捆绑的旧版 Eclipse 存在缺陷, 或者与 Scala IDE 4.4.1 捆绑的 Scala 版本存在版本不兼容问题。
我们建议改用 IntelliJ(见上文)
不要犹豫,尽管提问!
如果你需要任何帮助,请通过邮件列表联系开发者和社区。
如果你在 Flink 中发现了 bug,请提交一个 issue。
Apache Flink 的文档位于网站:https://flink.apache.org
或源码的 docs/ 目录中。
这是一个活跃的开源项目。我们始终欢迎想要使用该系统或为其做出贡献的人。 如果你正在寻找适合你技能的实现任务,请联系我们。 这篇文章介绍了如何为 Apache Flink 做贡献。
Apache Flink 是 Apache 软件基金会(ASF)的一个开源项目。 Apache Flink 项目起源于 Stratosphere 研究项目。