Flink入门案例

wen java案例 1

本文目录导读:

Flink入门案例

  1. 开发环境准备
  2. 完整WordCount案例
  3. 更完整的入门案例:用户行为分析
  4. 测试数据文件
  5. 运行方式
  6. 常见问题解决方案

我来为您提供一个完整的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,最后尝试更复杂的用户行为分析案例。

抱歉,评论功能暂时关闭!