本文目录导读:

我来为您提供一个完整的Flink入门案例,包含WordCount(词频统计)这一经典示例,以及详细的开发环境搭建和代码实现。
开发环境准备
Maven依赖配置
<properties>
<flink.version>1.17.1</flink.version>
<scala.version>2.12</scala.version>
</properties>
<dependencies>
<!-- Flink核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink流处理依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink客户端(用于本地调试) -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
</dependency>
</dependencies>
完整WordCount案例
1 批处理WordCount(DataSet API)
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.tuple.Tuple2;
import org.apache.flink.util.Collector;
public class BatchWordCount {
public static void main(String[] args) throws Exception {
// 1. 创建批处理执行环境
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// 2. 读取数据源
String inputPath = "input/words.txt";
if (args.length > 0) {
inputPath = args[0];
}
// 3. 数据处理流程
env.readTextFile(inputPath)
.flatMap(new FlatMapFunction<String, String>() {
@Override
public void flatMap(String line, Collector<String> out) {
// 按空格切分并输出每个单词
String[] words = line.toLowerCase().split("\\W+");
for (String word : words) {
if (word.length() > 0) {
out.collect(word);
}
}
}
})
.map(new MapFunction<String, Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> map(String word) {
return new Tuple2<>(word, 1);
}
})
.groupBy(0)
.sum(1)
.print();
}
}
2 流处理WordCount(DataStream API)
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.functions.KeySelector;
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;
public class StreamingWordCount {
public static void main(String[] args) throws Exception {
// 1. 创建流处理执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 设置并行度(可选)
env.setParallelism(1);
// 2. 读取数据源(这里使用 socket 流模拟实时数据)
// 运行时需要在终端执行: nc -lk 9999
DataStream<String> textStream = env.socketTextStream("localhost", 9999);
// 3. 数据处理流程
SingleOutputStreamOperator<Tuple2<String, Integer>> result = textStream
.flatMap(new Tokenizer())
.keyBy(new KeySelector<Tuple2<String, Integer>, String>() {
@Override
public String getKey(Tuple2<String, Integer> value) {
return value.f0;
}
})
.sum(1);
// 4. 输出结果
result.print().name("word-count-output");
// 5. 执行程序
env.execute("Streaming Word Count");
}
// 自定义FlatMap函数
public static class Tokenizer implements FlatMapFunction<String, Tuple2<String, Integer>> {
@Override
public void flatMap(String line, Collector<Tuple2<String, Integer>> out) {
String[] words = line.toLowerCase().split("\\W+");
for (String word : words) {
if (word.length() > 0) {
out.collect(new Tuple2<>(word, 1));
}
}
}
}
}
3 Lambda表达式版本(Java 8+)
import org.apache.flink.api.common.typeinfo.Types;
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.util.Collector;
public class LambdaWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> textStream = env.socketTextStream("localhost", 9999);
textStream
.flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
String[] words = line.toLowerCase().split("\\W+");
for (String word : words) {
if (word.length() > 0) {
out.collect(new Tuple2<>(word, 1));
}
}
})
.returns(Types.TUPLE(Types.STRING, Types.INT))
.keyBy(0)
.sum(1)
.print();
env.execute("Lambda Word Count");
}
}
更完整的入门案例:用户行为分析
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;
public class UserBehaviorAnalysis {
public static void main(String[] args) throws Exception {
// 创建执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 模拟用户行为数据源(格式:用户ID,行为,时间戳)
DataStream<String> input = env.socketTextStream("localhost", 8888);
// 解析并转换数据
DataStream<Tuple3<String, String, Long>> userBehavior = input
.map(new MapFunction<String, Tuple3<String, String, Long>>() {
@Override
public Tuple3<String, String, Long> map(String line) throws Exception {
String[] fields = line.split(",");
return new Tuple3<>(
fields[0], // 用户ID
fields[1], // 行为类型
Long.parseLong(fields[2]) // 时间戳
);
}
});
// 设置事件时间和watermark
DataStream<Tuple3<String, String, Long>> withWatermark = userBehavior
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<Tuple3<String, String, Long>>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner(new SerializableTimestampAssigner<Tuple3<String, String, Long>>() {
@Override
public long extractTimestamp(Tuple3<String, String, Long> element, long recordTimestamp) {
return element.f2;
}
})
);
// 统计每5分钟的用户行为数量
withWatermark
.map(new MapFunction<Tuple3<String, String, Long>, Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> map(Tuple3<String, String, Long> value) {
return new Tuple2<>(value.f1, 1);
}
})
.keyBy(0)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.sum(1)
.print();
env.execute("User Behavior Analysis");
}
}
测试数据文件
创建 input/words.txt 文件:
Hello Flink Hello World
Flink is a streaming framework
Hello Flink is powerful
Streaming processing with Flink
运行方式
本地运行
直接在IDE中运行主类即可。
命令行提交(需要Flink集群)
# 打包 mvn clean package # 提交到Flink集群 flink run -c com.example.WordCount /path/to/your-jar-with-dependencies.jar
Socket输入测试(流处理)
# 启动socket服务 nc -lk 9999 # 然后在控制台输入单词,如: Hello Flink Hello World
常见问题解决方案
1 处理依赖冲突
<!-- 在pom.xml中添加 -->
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.4.1</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>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
2 处理类型擦除问题
// 使用returns方法指定类型 stream.map(...).returns(Types.TUPLE(Types.STRING, Types.INT));
这个入门案例涵盖了Flink的核心概念,包括:
- 数据源处理(文件、socket)
- 转换操作(map、flatMap)
- 分组聚合(keyBy、sum)
- 窗口计算(tumbling window)
- 事件时间和水印
建议从批处理WordCount开始,然后过渡到流处理WordCount,最后尝试更复杂的用户行为分析案例。