温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

flink框架如何做流处理

发布时间:2026-09-06 16:16:20 来源:亿速云 阅读:88 作者:小樊 栏目:软件技术

Apache Flink 是一个分布式、高性能、有状态的流处理框架,核心设计理念是:一切皆为流(批处理只是流的一个特例)。下面从概念 → 开发流程 → 核心机制 → 示例几个层面,系统讲一下 Flink 如何做流处理。


一、Flink 流处理的核心概念

1. 数据流(DataStream)

Flink 流处理的基本抽象是 DataStream

  • 无界流(Unbounded):实时数据(如 Kafka、Socket)
  • 有界流(Bounded):历史数据(如文件,Flink 把批当流处理)

2. 基本组成

一个 Flink 流处理程序通常包含:

数据源(Source) → 转换(Transformation) → 输出(Sink)
  • Source:Kafka、Socket、文件、CDC 等
  • Transformation:map / filter / keyBy / window 等
  • Sink:Kafka、MySQL、HDFS、ClickHouse 等

二、Flink 流处理开发流程(以 Java 为例)

1. 创建执行环境

StreamExecutionEnvironment env =
    StreamExecutionEnvironment.getExecutionEnvironment();

2. 定义数据源(Source)

DataStream<String> stream = env
    .socketTextStream("localhost", 9999);

或从 Kafka:

KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("localhost:9092")
    .setTopics("topic")
    .setGroupId("group")
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

DataStream<String> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka");

3. 数据转换(Transformation)

简单转换

DataStream<Integer> nums = stream
    .map(Integer::parseInt);

过滤

DataStream<Integer> filtered = nums
    .filter(x -> x > 10);

分组 + 聚合(Keyed Stream)

DataStream<Tuple2<String, Integer>> result = stream
    .map(s -> Tuple2.of(s, 1))
    .keyBy(value -> value.f0)
    .sum(1);

4. 时间语义 & 窗口(流处理重点)

Flink 支持三种时间:

  • Event Time(事件时间,最常用)
  • Ingestion Time
  • Processing Time

设置事件时间 + 水位线

WatermarkStrategy<Event> strategy = WatermarkStrategy
    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((event, ts) -> event.timestamp);

窗口计算

stream
    .keyBy(e -> e.key)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .sum("value");

窗口类型:

  • 滚动窗口(Tumbling)
  • 滑动窗口(Sliding)
  • 会话窗口(Session)

5. 数据输出(Sink)

result.print(); // 控制台

或写入 Kafka / MySQL:

stream.addSink(new FlinkKafkaProducer<>(
    "localhost:9092",
    "output-topic",
    new SimpleStringSchema()
));

6. 提交执行

env.execute("Flink Stream Job");

三、Flink 流处理的关键机制

1. 状态管理(State)

流处理必须有状态:

  • ValueState
  • ListState
  • MapState

示例:

ValueState<Integer> countState = getRuntimeContext()
    .getState(new ValueStateDescriptor<>("count", Integer.class));

2. 精确一次语义(Exactly-Once)

通过:

  • Checkpoint
  • 两阶段提交(Kafka Sink)

保证数据不丢、不重。


3. 容错机制

  • Checkpoint:定期保存状态
  • Savepoint:手动保存,用于升级
  • 任务失败自动恢复

四、一个完整示例(WordCount 流处理)

StreamExecutionEnvironment env =
    StreamExecutionEnvironment.getExecutionEnvironment();

DataStream<String> text = env.socketTextStream("localhost", 9999);

text.flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
        for (String word : line.split(" ")) {
            out.collect(Tuple2.of(word, 1));
        }
    })
    .keyBy(value -> value.f0)
    .sum(1)
    .print();

env.execute("Socket WordCount");

五、Flink 流处理适用场景

  • 实时数仓
  • 实时监控 / 告警
  • 实时 ETL
  • 风控 / 反欺诈
  • CDC 数据同步

如果你愿意,我可以:

  • Scala / Python 再写一遍
  • Flink SQL 流处理
  • 对比 Spark Streaming vs Flink
  • 给你一个 Kafka → Flink → MySQL 实战方案

你想从哪一块深入?

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI