Flink入门宝典(详细截图版)


本文基于java构建Flink1.9版本入门程序,需要Maven 3.0.4 和 Java 8 以上版本。需要安装Netcat进行简单调试。

这里简述安装过程,并使用IDEA进行开发一个简单流处理程序,本地调试或者提交到Flink上运行,Maven与JDK安装这里不做说明。

一、Flink简介

Flink诞生于欧洲的一个大数据研究项目StratoSphere。该项目是柏林工业大学的一个研究性项目。早期,Flink是做Batch计算的,但是在2014年,StratoSphere里面的核心成员孵化出Flink,同年将Flink捐赠Apache,并在后来成为Apache的顶级大数据项目,同时Flink计算的主流方向被定位为Streaming,即用流式计算来做所有大数据的计算,这就是Flink技术诞生的背景。

2015开始阿里开始介入flink 负责对资源调度和流式sql的优化,成立了阿里内部版本blink在最近更新的1.9版本中,blink开始合并入flink,

未来flink也将支持java,scala,python等更多语言,并在机器学习领域施展拳脚。

二、Flink开发环境搭建

首先要想运行Flink,我们需要下载并解压Flink的二进制包,下载地址如下:https://flink.apache.org/downloads.html

我们可以选择Flink与Scala结合版本,这里我们选择最新的1.9版本Apache Flink 1.9.0 for Scala 2.12进行下载。

Flink在Windows和Linux下的安装与部署可以查看 Flink快速入门--安装与示例运行,这里演示windows版。

安装成功后,启动cmd命令行窗口,进入flink文件夹,运行bin目录下的start-cluster.bat

$ cd flink
$ cd bin
$ start-cluster.bat
Starting a local cluster with one JobManager process and one TaskManager process.
You can terminate the processes via CTRL-C in the spawned shell windows.
Web interface by default on http://localhost:8081/.

显示启动成功后,我们在浏览器访问 http://localhost:8081/可以看到flink的管理页面。

三、Flink快速体验

请保证安装好了flink,还需要Maven 3.0.4 和 Java 8 以上版本。这里简述Maven构建过程。

其他详细构建方法欢迎查看:快速构建第一个Flink工程

1、搭建Maven工程

使用Flink Maven Archetype构建一个工程。

 $ mvn archetype:generate                                     -DarchetypeGroupId=org.apache.flink                    -DarchetypeArtifactId=flink-quickstart-java            -DarchetypeVersion=1.9.0

你可以编辑自己的artifactId groupId

目录结构如下:

$ tree quickstart/
quickstart/
├── pom.xml
└── src
    └── main
        ├── java
        │   └── org
        │       └── myorg
        │           └── quickstart
        │               ├── BatchJob.java
        │               └── StreamingJob.java
        └── resources
            └── log4j.properties

在pom中核心依赖:

<dependencies>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-java</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java_2.11</artifactId>
        <version>${flink.version}</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-clients_2.11</artifactId>
        <version>${flink.version}</version>
    </dependency>
</dependencies>

2、编写代码

StreamingJob

import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.util.Collector;
public class StreamingJob {

    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        DataStream<Tuple2<String, Integer>> dataStreaming = env
                .socketTextStream("localhost", 9999)
                .flatMap(new Splitter())
                .keyBy(0)
                .timeWindow(Time.seconds(5))
                .sum(1);

        dataStreaming.print();

        // execute program
        env.execute("Flink Streaming Java API Skeleton");
    }
    public static class Splitter implements FlatMapFunction<String, Tuple2<String, Integer>> {

        @Override
        public void flatMap(String sentence, Collector<Tuple2<String, Integer>> out) throws Exception {
            for(String word : sentence.split(" ")){
                out.collect(new Tuple2<String, Integer>(word, 1));
            }
        }

    }
}

3、调试程序

安装netcat工具进行简单调试。

启动netcat 输入:

nc -l 9999

启动程序

在netcat中输入几个单词 逗号分隔

在程序一端查看结果

4、程序提交到Flink

启动flink

windows为 start-cluster.bat    linux为start-cluster.sh

localhost:8081查看管理页面

通过maven对代码打包

将打好的包提交到flink上

查看log

tail -f log/flink-***-jobmanager.out

在netcat中继续输入单词,在Running Jobs中查看作业状态,在log中查看输出。

四、Flink 编程模型

Flink提供不同级别的抽象来开发流/批处理应用程序。

最低级抽象只提供有状态流

在实践中,大多数应用程序不需要上述低级抽象,而是针对Core API编程,?如DataStream API(有界/无界流)和DataSet API(有界数据集)。

Table Api声明了一个表,遵循关系模型。

最高级抽象是SQL

我们这里只用到了DataStream API。

Flink程序的基本构建块是转换

一个程序的基本构成:

l?获取execution environment

l?加载/创建原始数据

l?指定这些数据的转化方法

l?指定计算结果的存放位置

l?触发程序执行

五、DataStreaming API使用

1、获取execution environment

StreamExecutionEnvironment是所有Flink程序的基础,获取方法有:

getExecutionEnvironment()

createLocalEnvironment()

createRemoteEnvironment(String host, int port, String ... jarFiles)

一般情况下使用getExecutionEnvironment。如果你在IDE或者常规java程序中执行可以通过createLocalEnvironment创建基于本地机器的StreamExecutionEnvironment。如果你已经创建jar程序希望通过invoke方式获取里面的getExecutionEnvironment方法可以使用createRemoteEnvironment方式。

2、加载/创建原始数据

StreamExecutionEnvironment提供的一些访问数据源的接口

(1)基于文件的数据源

readTextFile(path)
readFile(fileInputFormat, path)
readFile(fileInputFormat, path, watchType, interval, pathFilter, typeInfo)

(2)基于Socket的数据源(本文使用的)

l?socketTextStream

?

(3)基于Collection的数据源

fromCollection(Collection)
fromCollection(Iterator, Class)
fromElements(T ...)
fromParallelCollection(SplittableIterator, Class)
generateSequence(from, to)

3、转化方法

(1)Map方式:DataStream -> DataStream

功能:拿到一个element并输出一个element,类似Hive中的UDF函数

举例:

DataStream<Integer> dataStream = //...
dataStream.map(new MapFunction<Integer, Integer>() {
[email protected]
????public Integer map(Integer value) throws Exception {
????????return 2 * value;
????}
});

(2)FlatMap方式:DataStream -> DataStream

功能:拿到一个element,输出多个值,类似Hive中的UDTF函数

举例:

dataStream.flatMap(new FlatMapFunction<String, String>() {
[email protected]
????public void flatMap(String value, Collector<String> out)
????????throws Exception {
????????for(String word: value.split(" ")){
????????????out.collect(word);
????????}
????}
});

(3)Filter方式:DataStream -> DataStream

功能:针对每个element判断函数是否返回true,最后只保留返回true的element

举例:

dataStream.filter(new FilterFunction<Integer>() {
[email protected]
????public boolean filter(Integer value) throws Exception {
????????return value != 0;
????}
});

(4)KeyBy方式:DataStream -> KeyedStream

功能:逻辑上将流分割成不相交的分区,每个分区都是相同key的元素

举例:

dataStream.keyBy("someKey") // Key by field "someKey"
dataStream.keyBy(0) // Key by the first element of a Tuple

(5)Reduce方式:KeyedStream -> DataStream

功能:在keyed data stream中进行轮训reduce。

举例:

keyedStream.reduce(new ReduceFunction<Integer>() {
[email protected]
????public Integer reduce(Integer value1, Integer value2)
????throws Exception {
????????return value1 + value2;
????}
});

(6)Aggregations方式:KeyedStream -> DataStream

功能:在keyed data stream中进行聚合操作

举例:

keyedStream.sum(0);
keyedStream.sum("key");
keyedStream.min(0);
keyedStream.min("key");
keyedStream.max(0);
keyedStream.max("key");
keyedStream.minBy(0);
keyedStream.minBy("key");
keyedStream.maxBy(0);
keyedStream.maxBy("key");

(7)Window方式:KeyedStream -> WindowedStream

功能:在KeyedStream中进行使用,根据某个特征针对每个key用windows进行分组。

举例:

dataStream.keyBy(0).window(TumblingEventTimeWindows.of(Time.seconds(5))); // Last 5 seconds of data

(8)WindowAll方式:DataStream -> AllWindowedStream

功能:在DataStream中根据某个特征进行分组。

举例:

dataStream.windowAll(TumblingEventTimeWindows.of(Time.seconds(5))); // Last 5 seconds of data

(9)Union方式:DataStream* -> DataStream

功能:合并多个数据流成一个新的数据流

举例:

dataStream.union(otherStream1, otherStream2, ...);

(10)Split方式:DataStream -> SplitStream

功能:将流分割成多个流

举例:

SplitStream<Integer> split = someDataStream.split(new OutputSelector<Integer>() {
[email protected]
????public Iterable<String> select(Integer value) {
????????List<String> output = new ArrayList<String>();
????????if (value % 2 == 0) {
????????????output.add("even");
????????}
????????else {
????????????output.add("odd");
????????}
????????return output;
????}
});

(11)Select方式:SplitStream -> DataStream

功能:从split stream中选择一个流

举例:

SplitStream<Integer> split;
DataStream<Integer> even = split.select("even");
DataStream<Integer> odd = split.select("odd");
DataStream<Integer> all = split.select("even","odd");

4、输出数据

writeAsText()
writeAsCsv(...)
print() / printToErr()
writeUsingOutputFormat() / FileOutputFormat
writeToSocket
addSink

更多Flink相关原理:

穿梭时空的实时计算框架——Flink对时间的处理

大数据实时处理的王者-Flink

统一批处理流处理——Flink批流一体实现原理

Flink快速入门--安装与示例运行

快速构建第一个Flink工程

更多实时计算,Flink,Kafka等相关技术博文,欢迎关注实时流式计算:

原文地址:https://www.cnblogs.com/tree1123/p/11539955.html

时间: 2024-07-31 19:39:47

Flink入门宝典(详细截图版)的相关文章

React Native 入门宝典

声明:该书的笔者为徐嬴老师,一名具有5年IOS开发经验,和两年RN开发经验的老司机. 原文可以在gitbook上找到 笔者只是为他的书中提的的一些列问题,进行有偿答疑. 有偿答疑.本书将持续保持更新,有关问题可以加群讨论. 正在上传...取消 简介 笔者在研究ReactNative过程中,发现其中文资料相对较少,已出版的大部分图书资料都已过时.Facebook中的ReactNative开发团队以每月更新一版的速度在向前推进版本. 为更好的让广大开发者快速入门ReactNative,笔者结合自身开

c语言入门经典(第5版)

文章转载:http://mrcaoyc.blog.163.com/blog/static/23939201520159135915734 文件大小:126MB 文件格式:PDF    [点击下载] C语言入门经典(第5版)  内容简介: C语言是每一位程序员都应该掌握的基础语言.C语言是微软.NET编程中使用的C#语言的基础:C语言是iPhone.iPad和其他苹果设备编程中使用的Objective-C语言的基础:C语言是在很多环境中(包括GNU项目)被广泛使用的C++语言的基础.C语言也是Li

Expo 入门宝典 一 (Quick Start)

本人决心翻译Expo,为学习Rn(react native)的学习者提供帮助.传统上Rn开发,优势都在Mac Ios ,很少有用Windows andriod开发的,而2017年上线的Expo为我们广大windows做Rn开发提供了很大的便利条件.Rn开发也迎来了春天. 关于Rn的简单说明,目前市场上主流的两大移动端系统,Android 和 Ios,而开发这两个系统上的App,传统上,分为Ios开发和Android开发,这就有一个问题,一个公司要上线一款app,但是需要至少需要一个Ios开发,和

powershell入门教程-v0.3版

powershell入门教程-v0.3版 来源 https://www.itsvse.com/thread-3650-1-1.html 参考 http://www.cnblogs.com/piapia/ https://www.pstips.net/powershell-online-tutorials http://www.cnblogs.com/volcanol/tag/PowerShell/ 问:如何开启powershell脚本运行权限?答:echo 下面代码,在管理员权限cmd中运行,在

Flink入门(五)——DataSet Api编程指南

Apache Flink Apache Flink 是一个兼顾高吞吐.低延迟.高性能的分布式处理框架.在实时计算崛起的今天,Flink正在飞速发展.由于性能的优势和兼顾批处理,流处理的特性,Flink可能正在颠覆整个大数据的生态. DataSet API 首先要想运行Flink,我们需要下载并解压Flink的二进制包,下载地址如下:https://flink.apache.org/downloads.html 我们可以选择Flink与Scala结合版本,这里我们选择最新的1.9版本Apache

Tungsten Fabric入门宝典丨TF组件的七种“武器”

Tungsten Fabric入门宝典系列文章,来自技术大牛倾囊相授的实践经验,由TF中文社区为您编译呈现,旨在帮助新手深入理解TF的运行.安装.集成.调试等全流程.如果您有相关经验或疑问,欢迎与我们互动,并与社区极客们进一步交流.更多TF技术文章,请点击公号底部按钮>学习>文章合集. 作者:Tatsuya Naganawa 译者:TF编译组 Tungsten Fabric中有很多不同的组件.接下来我简要描述它们的用法. 概览 总体而言,Tungsten Fabric中包含7种角色和(多达)3

在intellij IDEA中为web应用创建图片虚拟目录(详细截图)

在intellij IDEA中为web应用创建图片虚拟目录(详细截图) 在intellij IDEA中为web应用创建图片虚拟目录详细截图 工程配置和环境 操作步骤 在非IDE环境下配置虚拟目录 本文主要展示如何在intellij IDEA中为web应用添加虚拟目录映射,并附上步骤截图 工程配置和环境 我使用的版本为 tomcat 8.0.30 intellij 15.0.2 jdk 1.8.0_25 已经部署好了一个web应用,并且已经在IDEA中添加好了tomcat容器,现在想为这个web应

2015年最新Android基础入门教程目录(完结版)

2015年最新Android基础入门教程目录(完结版) 标签(空格分隔): Android基础入门教程 前言: 关于<2015年最新Android基础入门教程目录>终于在今天落下了帷幕,全套教程 共148节已编写完毕,附上目录,关于教程的由来,笔者的情况和自学心得,资源分享 以及一些疑问等可戳:<2015最新Android基础入门教程>完结散花~ 下面是本系列教程的完整目录: 第一章:环境搭建与开发相关(已完结 10/10) Android基础入门教程--1.1 背景相关与系统架构

CENTOS7 安装openstack mitaka版本(最新整理完整版附详细截图和操作步骤,添加了cinder和vxlan)

CENTOS7 安装openstack mitaka版本(最新整理完整版附详细截图和操作步骤,添加了cinder和vxlan,附上个节点的配置文件) 实验环境准备: 为了更好的实现分布式mitaka版本的效果.我才有的是VMware的workstations来安装三台虚拟机,分别来模拟openstack的controller节点 compute节点和cinder节点.(我的宿主机配置为 500g 硬盘 16g内存,i5cpu.强烈建议由条件的朋友将内存配置大一点,因为我之前分配的2g太卡.) 注