Apache Flink 是一个分布式、高性能、有状态的流处理框架,核心设计理念是:一切皆为流(批处理只是流的一个特例)。下面从概念 → 开发流程 → 核心机制 → 示例几个层面,系统讲一下 Flink 如何做流处理。
Flink 流处理的基本抽象是 DataStream:
一个 Flink 流处理程序通常包含:
数据源(Source) → 转换(Transformation) → 输出(Sink)
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
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");
DataStream<Integer> nums = stream
.map(Integer::parseInt);
DataStream<Integer> filtered = nums
.filter(x -> x > 10);
DataStream<Tuple2<String, Integer>> result = stream
.map(s -> Tuple2.of(s, 1))
.keyBy(value -> value.f0)
.sum(1);
Flink 支持三种时间:
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");
窗口类型:
result.print(); // 控制台
或写入 Kafka / MySQL:
stream.addSink(new FlinkKafkaProducer<>(
"localhost:9092",
"output-topic",
new SimpleStringSchema()
));
env.execute("Flink Stream Job");
流处理必须有状态:
示例:
ValueState<Integer> countState = getRuntimeContext()
.getState(new ValueStateDescriptor<>("count", Integer.class));
通过:
保证数据不丢、不重。
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");
如果你愿意,我可以:
你想从哪一块深入?
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。