Java Flink案例

wen java案例 1

本文目录导读:

Java Flink案例

  1. 环境准备
  2. 案例一:实时流处理 WordCount
  3. 案例二:批处理 WordCount
  4. 案例三:带事件时间与窗口的流处理
  5. 运行与测试
  6. 分布式集群运行(可选)
  7. 关键点总结

下面我给你一个 完整的 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 数据源,如 socketTextStreamreadTextFile
Transformation 转换操作,如 flatMapmapkeyBysum
Sink 数据输出,如 printwriteAsText
Window 窗口操作,处理无界流中的有界数据
Watermark 用于处理事件时间的乱序问题

如果你需要更复杂的案例(如 连接 Kafka使用 Table APICEP 复杂事件处理状态管理 等),请告诉我,我可以为你补充。

上一篇DataFrame案例

下一篇Spark SQL案例

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