You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink反序列化阶段如何将单条Kinesis记录拆分为多条写入Elasticsearch

核心实现逻辑

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 23:27:00