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

Apache Flink已废弃SplitStream如何使用side-outputs替代?

Side Outputs 替换 SplitStream 实现方案

SplitStream 是 Flink 早期提供的流拆分 API,目前已完全废弃,侧输出流(Side Outputs)是官方唯一推荐的替代方案,具体替换步骤如下:


1. 定义分流标识 OutputTag

对应旧 API 中拆分用的字符串标签,你需要为每一个拆分后的子流定义唯一的 OutputTag 实例,实例的泛型和你要拆分的子流数据类型保持一致:

// 示例对应原标签"some-string",替换为你实际的业务标签和数据类型
OutputTag<你的数据类型> targetTag = OutputTag.forName(
    "some-string",
    TypeInformation.of(new TypeHint<你的数据类型>() {})
);

2. 主流处理中分发数据到侧输出

用 ProcessFunction 处理原始主数据流,按照你的业务分流规则,将符合条件的数据发送到对应 OutputTag 标识的侧输出流中:

SingleOutputStreamOperator<主流原始数据类型> processedMainStream = mainDataStream.process(
    new ProcessFunction<主流原始数据类型, 主流原始数据类型>() {
        @Override
        public void processElement(
            主流原始数据类型 value,
            Context ctx,
            Collector<主流原始数据类型> out
        ) throws Exception {
            // 此处替换为你实际的分流判断逻辑
            if (数据符合"some-string"标签的筛选条件) {
                // 发送到侧输出流
                ctx.output(targetTag, value);
            }
            // 不需要保留主流数据的话可以删除下面这行
            out.collect(value);
        }
    }
);

3. 提取目标侧输出流

这一步完全等价于你原代码中的 mainDataStream.select("some-string") 逻辑,直接从处理后的主流中提取对应标签的侧输出流即可:

// 得到的targetStream就等价于原SplitStream.select后得到的子流
DataStream<你的数据类型> targetStream = processedMainStream.getSideOutput(targetTag);

相比旧的SplitStream,侧输出流支持不同子流使用不同数据类型,也支持单条数据匹配多个分流规则,灵活性和性能都更优。

内容的提问来源于stack exchange,提问作者user4202236

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 10:06:02