Flink反序列化阶段如何将单条Kinesis记录拆分为多条写入Elasticsearch
Flink 单条Kinesis记录拆分为多条记录实现方案
核心实现逻辑
Flink 内置的FlatMapFunction算子天生支持单输入元素输出多元素的一对多转换场景,完全匹配你的需求,只需在现有的反序列化逻辑之后追加该算子即可。
具体实现步骤
1. 前置准备(你已完成可跳过)
已定义两类POJO:
- 整包事件结构:示例命名为
KinesisLogPackage,包含logEvents列表属性 - 单条日志结构:示例命名为
LogEvent,对应logEvents子项的字段结构
2. 实现FlatMap转换逻辑
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.util.Collector; // 输入为整包POJO,输出为单条日志POJO public class LogEventSplitter implements FlatMapFunction<KinesisLogPackage, LogEvent> { @Override public void flatMap(KinesisLogPackage input, Collector<LogEvent> out) throws Exception { // 先判空避免空指针异常 if (input != null && input.getLogEvents() != null) { // 遍历列表逐条输出 for (LogEvent event : input.getLogEvents()) { // 可选:如果需要把整包的logGroup、logStream等公共字段带到单条日志中,可在这里补充赋值 // event.setLogGroup(input.getLogGroup()); // event.setLogStream(input.getLogStream()); out.collect(event); } } } }
3. 串联到现有流处理流程
将拆分算子插入到反序列化之后、Elasticsearch写入之前即可:
// 现有逻辑:Kinesis数据源 -> 解压 -> 反序列化为整包POJO DataStream<KinesisLogPackage> packageStream = env .addSource(new FlinkKinesisConsumer<>(...)) .map(new GzipDecompressFunction()) .map(str -> objectMapper.readValue(str, KinesisLogPackage.class)); // 新增拆分逻辑,得到单条日志流 DataStream<LogEvent> singleEventStream = packageStream.flatMap(new LogEventSplitter()); // 现有逻辑:直接写入Elasticsearch即可 singleEventStream.addSink(new ElasticsearchSink<>(...));
扩展说明
如果你的后续写入逻辑需要保留原始Kinesis记录的元数据(比如分片ID、摄入时间等),也可以改用ProcessFunction实现拆分,在ProcessFunction中可以通过上下文获取到原始记录的元数据信息,补充到单条日志中。
内容的提问来源于stack exchange,提问作者Seungjin Lee
相关产品推荐
相关产品推荐

