image-20240627224250581

  • 最顶层:SQL/Table API 提供了操作关系表、执行SQL语句分析的API库,供我们方便的开发SQL相关程序
  • 中层:流和批处理API层,提供了一系列流和批处理的API和算子供我们对数据进行处理和分析
  • 最底层:运行时层,提供了对Flink底层关键技术的操纵,如对Event、state、time、window等进行精细化控制的操作API

DataStreamAPI

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.operators.*;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;

public class WordCountBatchDataStream {
public static void main(String[] args) throws Exception {
/**
* 实现步骤:
* 1、初始化flink批处理运行环境
* 2、从指定文件中读取数据
* 3、对获取的数据按指定分隔符切分
* 4、对切分后的每个单词计数1
* 5、对相同单词的进行分组操作
* 6、对分组后的单词进行累加操作
* 7、打印输出
* 8、启动作业、提交任务
*/
// 1, 初始化Spark运行环境
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// 2, 从指定文件中读取数据
DataSource<String> text = env.readTextFile("D:\\Download\\wordcount.txt");
// 3, 对获取的数据按指定分隔符切分
FlatMapOperator<String, String> wordDS = text.flatMap(new FlatMapFunction<String, String>() {
/**
* 按照逗号进行切割
* @param value 输入的字符串
* @param out 返回结果的收集器
* @throws Exception
*/
@Override
public void flatMap(String value, Collector<String> out) throws Exception {
// 将输入的value按照','进行切割
String[] strings = value.split(",");
// 使用增强for循环遍历数组(快捷键 foreach
for (String str :
strings) {
out.collect(str);
}
}
});
// 4, 对切分后的每个单词计数 1, 注意Flink提供了一种类似于map的数据形式Tuple2, 二元组形式.
MapOperator<String, Tuple2<String,Integer>> mapDs = wordDS.map(new MapFunction<String, Tuple2<String,Integer>>() {
/**
*
* @param value 输入的字符串
* @return (word, 1)
* @throws Exception
*/
@Override
public Tuple2<String,Integer> map(String value) throws Exception {
return Tuple2.of(value, 1);
}
});
// 5, 对相同的单词进行分组操作
UnsortedGrouping<Tuple2<String,Integer>> groupBy = mapDs.groupBy(0);
// 6, 对分组后的单词进行累加操作
AggregateOperator<Tuple2<String,Integer>> sumDs = groupBy.sum(1);

// 7, 打印输出
sumDs.print();
}
}

TableAPI

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
import org.apache.flink.table.api.*;


import javax.swing.text.TableView;

import static org.apache.flink.table.api.Expressions.$;

public class WordCountBatchTable {
public static void main(String[] args) {
/**
* 1、构建flink table批处理运行环境
* 2、创建Source Table获取数据
* 3、对获取到的数据进行计算
* 4、创建Sink Table用于输出数据
* 5、将计算结果写入到Sink Table
*/
// 1, 构建Flink批处理运行环境
EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build();
TableEnvironment tEnv = TableEnvironment.create(settings);

// 2, 创建Source Table获取数据
tEnv.createTemporaryTable("sourceTable", TableDescriptor.forConnector("filesystem")
.schema(Schema.newBuilder()
.column("userId", DataTypes.STRING())
.column("timestamp",DataTypes.BIGINT())
.column("money", DataTypes.DOUBLE())
.column("category", DataTypes.STRING())
.build())
.option("path", "D:\\Download\\order(1).csv")
.option("format", "csv")
.build());
// 3, 对获取到的数据进行计算
Table result = tEnv.from("sourceTable")
.groupBy($("userId"))
.select($("userId"), $("money").sum().as("totalMoney"));

//4, 创建Sink Table用于输出数据
// 结果表与输出表的schema必须一致(列的数量和列的数据类型)
tEnv.createTemporaryTable("sinkTable", TableDescriptor.forConnector("print")
.schema(Schema.newBuilder()
.column("userId", DataTypes.STRING())
.column("totalMoney", DataTypes.DOUBLE())
.build())
.build());
// 5, 将计算结果写入到SinkTable
result.executeInsert("sinkTable");
}
}

报错解决

Exception in thread “main” java.lang.reflect.InaccessibleObjectException: Unable to make field private static final long java.util.LinkedHashMap.serialVersionUID accessible: module java.base does not “opens java.util” to unnamed module @424ebba3

image-20240627223310715

image-20240627223417992