本文目录导读:

下面我给你一个 完整的 Apache Flink 入门案例,使用 Java 语言实现。
这个案例是经典的 “实时单词计数”(WordCount),包含流处理和批处理两种方式。
环境准备
1 Maven 依赖 (pom.xml)
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>flink-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<packaging>jar</packaging>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<flink.version>1.17.1</flink.version>
<scala.binary.version>2.12</scala.binary.version>
</properties>
<dependencies>
<!-- Flink Core API -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink Streaming API -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink Client (用于本地调试) -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- 日志框架 -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
<version>1.7.36</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>1.7.36</version>
</dependency>
</dependencies>
<build>
<plugins>
<!-- 打包插件 -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.4</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.example.StreamingWordCount</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
案例一:实时流处理 WordCount
从 Socket (TCP) 实时读取文本,进行单词计数。
package com.example;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
/**
* 实时流处理 WordCount 案例
*
* 启动方式:
* 1. 先启动 netcat: nc -lk 9999
* 2. 然后运行此程序
* 3. 在 netcat 窗口输入文本,观察控制台输出
*/
public class StreamingWordCount {
public static void main(String[] args) throws Exception {
// 1. 创建流处理执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 设置并行度
env.setParallelism(1);
// 3. 从 Socket 读取数据源 (nc -lk 9999)
DataStream<String> textStream = env.socketTextStream("localhost", 9999);
// 4. 数据处理管道
SingleOutputStreamOperator<Tuple2<String, Integer>> wordCountStream = textStream
// 4.1 数据分割与扁平化
.flatMap(new Tokenizer())
// 4.2 给每个单词计数 1
.map(new MapFunction<String, Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> map(String word) throws Exception {
return new Tuple2<>(word, 1);
}
})
// 4.3 按单词分组
.keyBy(0)
// 4.4 求和
.sum(1);
// 5. 输出结果到控制台
wordCountStream.print();
// 6. 执行任务
env.execute("Streaming Word Count");
}
/**
* 分词器:将一行文本拆分为单词
*/
public static class Tokenizer implements FlatMapFunction<String, String> {
@Override
public void flatMap(String line, Collector<String> out) throws Exception {
// 将非字母字符替换为空格,全部转小写,然后按空格分割
String[] words = line.toLowerCase()
.replaceAll("[^a-zA-Z\\s]", " ")
.trim()
.split("\\s+");
for (String word : words) {
if (word.length() > 0) {
out.collect(word);
}
}
}
}
}
案例二:批处理 WordCount
处理静态文件,一次性输出结果。
package com.example;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;
/**
* 批处理 WordCount 案例
*
* 输入: 一个文本文件
* 输出: 每个单词出现的次数(按次数降序)
*/
public class BatchWordCount {
public static void main(String[] args) throws Exception {
// 1. 创建批处理执行环境(Flink 1.17 中已标记为 deprecated)
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// 2. 从文件读取数据
String inputPath = "input/words.txt"; // 可以修改为实际文件路径
DataSet<String> text = env.readTextFile(inputPath);
// 3. 数据处理
DataSet<Tuple2<String, Integer>> wordCounts = text
// 扁平化:每行拆成单词
.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() {
@Override
public void flatMap(String line, Collector<Tuple2<String, Integer>> out) {
String[] words = line.toLowerCase().split("\\s+");
for (String word : words) {
out.collect(new Tuple2<>(word, 1));
}
}
})
// 按单词分组
.groupBy(0)
// 求和
.sum(1);
// 4. 按次数降序排序(可选)
DataSet<Tuple2<String, Integer>> sorted = wordCounts
.sortPartition(1, org.apache.flink.api.common.operators.Order.DESCENDING)
.setParallelism(1);
// 5. 输出到控制台
sorted.print();
// 6. 可以写入文件(可选)
// sorted.writeAsText("output/result.txt");
// 7. 执行任务(批处理会自动执行,但显式调用更规范)
// env.execute("Batch Word Count");
}
}
案例三:带事件时间与窗口的流处理
这是一个更高级的案例,演示 事件时间、Watermark 和 滚动窗口 的使用。
package com.example;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import java.time.Duration;
/**
* 带事件时间与滚动窗口的流处理
*
* 数据格式: timestamp,word,count
* 1000,hello,2
* 2000,world,3
*/
public class WindowingWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// 从 Socket 读取数据:格式为 timestamp,word,count
DataStream<String> text = env.socketTextStream("localhost", 9999);
// 解析为 (word, count, timestamp) 三元组
DataStream<Tuple3<String, Integer, Long>> parsed = text
.map(new MapFunction<String, Tuple3<String, Integer, Long>>() {
@Override
public Tuple3<String, Integer, Long> map(String line) throws Exception {
String[] parts = line.split(",");
return new Tuple3<>(
parts[1], // word
Integer.parseInt(parts[2]), // count
Long.parseLong(parts[0]) // event time
);
}
});
// 设置 Watermark 策略:允许 5 秒乱序
WatermarkStrategy<Tuple3<String, Integer, Long>> watermarkStrategy =
WatermarkStrategy
.<Tuple3<String, Integer, Long>>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner(new SerializableTimestampAssigner<Tuple3<String, Integer, Long>>() {
@Override
public long extractTimestamp(Tuple3<String, Integer, Long> element, long recordTimestamp) {
return element.f2; // 事件时间
}
});
// 应用 Watermark
DataStream<Tuple3<String, Integer, Long>> withWatermark =
parsed.assignTimestampsAndWatermarks(watermarkStrategy);
// 窗口计算:每 10 秒一个滚动窗口
DataStream<Tuple2<String, Integer>> windowedCounts =
withWatermark
.keyBy(t -> t.f0) // 按单词分组
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.sum(1); // 对 count 字段求和
// 输出
windowedCounts.print();
env.execute("Windowing Word Count");
}
}
运行与测试
1 运行流处理案例
# 1. 首先启动一个 TCP 服务 nc -lk 9999 # 2. 运行 Java 程序 # 方式一:IDE 中直接运行 StreamingWordCount # 方式二:打包后运行 mvn clean package java -jar target/flink-demo-1.0-SNAPSHOT.jar # 3. 在 netcat 窗口输入数据 # 然后观察控制台输出
2 运行批处理案例
# 准备输入文件 input/words.txt,内容例如: hello world flink java hello flink # 运行 BatchWordCount # 观察控制台输出
分布式集群运行(可选)
如果要在 Flink 集群上运行:
# 1. 启动 Flink 集群
$FLINK_HOME/bin/start-cluster.sh
# 2. 提交任务
$FLINK_HOME/bin/flink run \
-m yarn-cluster \
-c com.example.StreamingWordCount \
target/flink-demo-1.0-SNAPSHOT.jar
# 3. 查看 Web UI
# http://localhost:8081
关键点总结
| 概念 | 说明 |
|---|---|
| DataStream | 流式数据,无限数据集 |
| DataSet | 批处理数据集(1.17 后逐渐被 Table API 替代) |
| Source | 数据源,如 socketTextStream、readTextFile |
| Transformation | 转换操作,如 flatMap、map、keyBy、sum |
| Sink | 数据输出,如 print、writeAsText |
| Window | 窗口操作,处理无界流中的有界数据 |
| Watermark | 用于处理事件时间的乱序问题 |
如果你需要更复杂的案例(如 连接 Kafka、使用 Table API、CEP 复杂事件处理、状态管理 等),请告诉我,我可以为你补充。