摘要:一条实时数据流里总混着正常、异常、迟到、需监控四类数据,传统 filter 方案要遍历 N 遍、逻辑散落 N 个算子;Flink 侧输出流(Side Output)用一次遍历完成多路分发。这篇文章拆透侧输出:从 1+N 流模型、OutputTag 的类型信息机制(为什么必须带花括号)、ctx.output 到旁路记录通道的内部实现,到窗口迟到数据的专用 API(allowedLateness + sideOutputLateData)、双流对账异常旁路等完整代码案例,最后给出 side output vs filter vs union 的选型边界。
关键词:Flink 侧输出流、Side Output、OutputTag、ctx.output、getSideOutput、迟到数据、allowedLateness、sideOutputLateData、多流分发、旁路通道、TypeInformation、代码实现
一、一条流里混着四类数据,怎么办
实时业务里,一条输入流几乎总是混着:正常业务数据、格式异常/校验失败的数据、迟到数据(过了窗口触发时间才到)、需要监控的告警数据。经典做法是 N 个 filter 串联:
DataStream<Order> normal = orders.filter(o -> isValid(o));
DataStream<Order> invalid = orders.filter(o -> !isValid(o));
问题很明显:每条记录被遍历 N 次,N 个条件就是 N 倍处理开销;而且 filter 不命中的记录直接丢弃——想保留原始报文排障?做不到。侧输出流(Side Output)解决的就是这件事:一条流一次遍历,按条件拆成 1 条主流 + N 条旁路,各走各的下游,记录不丢。
二、侧输出流 1+N 模型

核心模型就一句话:在 process(...) 里,out.collect() 发射到主流,ctx.output(tag, value) 发射到 tag 对应的旁路。每个 OutputTag 定义一条独立旁路,旁路的数据类型可以完全不同——主流是 Order,旁路可以是 Alert、是 String 报文,互不约束。
三条核心特性,决定了它和 filter 的本质区别:
- 一次遍历,N 路分发:数据只被算子处理一遍,各旁路各取所需;
- 类型自由:每个 OutputTag 独立泛型,旁路可携带完全不同的数据类型(比如异常流直接带原始 JSON 报文);
- 状态级管道:旁路记录与主流共享 checkpoint 和 watermark——不丢数据、水位线语义一致。
还有一个容易被忽略的点:侧输出流本身也是普通 DataStream。getSideOutput(tag) 拿到的流可以继续 keyBy、开窗口、process、Sink,甚至可以再产生自己的侧输出(嵌套旁路)。
三、内部实现:OutputTag 为什么必须带花括号

3.1 类型信息:花括号的真相
侧输出流的类型信息来自 OutputTag 本身,而不是像主流那样由算子泛型推断。而 Java 的泛型是运行时擦除的——new OutputTag<>("late") 在运行时根本不知道元素类型是什么。解决办法是匿名子类:
// ✅ 匿名子类:子类继承时泛型被固化,TypeExtractor 能提取出 TypeInformation<String>
OutputTag<String> lateTag = new OutputTag<String>("late") {};
// ❌ 不带花括号:类型擦除后拿到的是 Object,序列化器无从构造
OutputTag<String> lateTag = new OutputTag<>("late");
不带花括号的版本在作业提交时往往不报错,运行期侧输出第一次反序列化时才炸——类型对不上,错误难排查。所以规范是:OutputTag 一律 new OutputTag<X>("id") {},并且定义成 static final(避免捕获外部 this 的序列化问题)。
3.2 发射与通道
ctx.output(tag, value) 的内部实现是 SideOutputDataOutput(Output 接口的实现),它按 tag 找到旁路通道,把记录写进侧输出分支。这里有两个容易误解的细节:
- 旁路与主流共享输出管线:记录走同一个 RecordWriter 输出缓冲区,所以 watermark 会照常推进到侧输出流的下游——迟到数据处理、窗口逻辑在旁路上依然生效;
- 侧输出本身不做重分区:记录往哪个并行子任务走,由下游算子的分区策略决定,
ctx.output不改变数据分布。
3.3 取流
下游 main.getSideOutput(tag) 按 tag 重建独立 DataStream——tag 需要和发射时是同一个对象(或 equals 相等),类型由 tag 携带的 TypeInformation 决定。这之后,它就是一个普通流,想怎么处理怎么处理。
四、代码实现:四个完整案例
4.1 实时数仓分流:正常 / 异常 / 延迟三路
ODS 层原始订单流进来,一次 process 拆成三路,异常流带原始报文、延迟流单独标记,全部不丢:
public class OrderSplitter extends ProcessFunction<Order, Order> {
private static final OutputTag<Order> INVALID_TAG = new OutputTag<Order>("invalid") {};
private static final OutputTag<Order> LATE_TAG = new OutputTag<Order>("late") {};
@Override
public void processElement(Order o, Context ctx, Collector<Order> out) {
// 1) 校验失败 → 异常旁路(保留原始数据供排障)
if (o.orderId == null || o.amount <= 0) {
ctx.output(INVALID_TAG, o);
return; // 校验失败的不进主流
}
// 2) 事件时间明显滞后于当前水位 → 延迟旁路
long wm = ctx.timerService().currentWatermark();
if (wm > Long.MIN_VALUE && o.eventTs < wm - 60_000L) {
ctx.output(LATE_TAG, o);
return;
}
// 3) 正常订单 → 主流(继续 DWD 加工)
out.collect(o);
}
}
DataStream<Order> main = orders.process(new OrderSplitter());
DataStream<Order> invalid = main.getSideOutput(OrderSplitter.INVALID_TAG);
DataStream<Order> late = main.getSideOutput(OrderSplitter.LATE_TAG);
// invalid → 对账/人工处理队列;late → 延迟补偿链路;main → 正常加工
注意 return 的用法:分派是互斥的,命中旁路就 return,避免一条记录同时进主流和旁路。这是侧输出写法里最常见的逻辑 bug——漏了 return,数据就"双写"了。
4.2 窗口迟到数据:allowedLateness + sideOutputLateData
这是侧输出最经典的生产场景。窗口触发计算后,allowedLateness 允许延迟内的迟到数据重新触发窗口(增量更新),而 allowedLateness 之外的迟到数据,通过 sideOutputLateData 进旁路——补算或修正,不污染已触发的结果:
OutputTag<Order> lateTag = new OutputTag<Order>("window-late") {};
DataStream<WindowResult> result = orders
.keyBy(Order::getProductId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.allowedLateness(Time.minutes(1)) // 窗口触发后再等 1 分钟
.sideOutputLateData(lateTag) // 1 分钟后还迟到的 → 旁路
.aggregate(new AmountAgg(), new WindowStat());
// 迟到旁路:单独做补算任务(T+1 修正 / 告警人工介入)
DataStream<Order> lateOrders = result.getSideOutput(lateTag);
两个关键参数要配套理解:
allowedLateness(1min):窗口触发后 1 分钟内到达的迟到数据,会重新触发窗口计算(结果增量更新);sideOutputLateData(tag):超过 allowedLateness 的迟到数据不再触发窗口,而是进旁路。
线上常见的错误是只配了 sideOutputLateData 忘配 allowedLateness——那样迟到数据既不会触发窗口、也进不了旁路,直接被静默丢弃,补数链路形同虚设。
4.3 双流对账异常旁路
订单流 + 支付流对账,匹配失败的记录进旁路(未匹配订单、未匹配支付分开),主流只保留匹配成功的结果:
public class Reconcile extends CoProcessFunction<Order, Pay, MatchResult> {
private static final OutputTag<Order> UNMATCHED_ORDER = new OutputTag<Order>("unmatched-order") {};
private static final OutputTag<Pay> UNMATCHED_PAY = new OutputTag<Pay>("unmatched-pay") {};
private ValueState<Order> pending;
@Override
public void processElement1(Order o, Context ctx, Collector<MatchResult> out) {
pending.update(o);
ctx.timerService().registerEventTimeTimer(o.ts + 10 * 60 * 1000L);
}
@Override
public void processElement2(Pay p, Context ctx, Collector<MatchResult> out) {
Order o = pending.value();
if (o == null) {
ctx.output(UNMATCHED_PAY, p); // 只有支付没有订单 → 异常旁路
} else {
pending.clear();
out.collect(new MatchResult(o, p, MATCHED));
}
}
@Override
public void onTimer(long ts, OnTimerContext ctx, Collector<MatchResult> out) {
Order o = pending.value();
if (o != null) {
ctx.output(UNMATCHED_ORDER, o); // 10 分钟没匹配 → 订单异常旁路
pending.clear();
}
}
}
注意两个侧输出 tag 的类型不同:UNMATCHED_ORDER 是 Order、UNMATCHED_PAY 是 Pay——侧输出流的类型自由度在这里直接体现,两个异常流可以各自接各自的修复链路。
4.4 定时器告警旁路
在 KeyedProcessFunction 里,定时器到期的告警走侧输出,业务流保持干净(监控逻辑不会因为告警 Sink 出问题而拖垮主流程):
private static final OutputTag<Alert> ALERT_TAG = new OutputTag<Alert>("alert") {};
@Override
public void onTimer(long ts, OnTimerContext ctx, Collector<Order> out) throws Exception {
Integer c = count.value();
if (c != null && c < expected) {
// 告警 → 侧输出流,单独接告警中心
ctx.output(ALERT_TAG, new Alert(ctx.getCurrentKey(), ts, "count-lag"));
}
count.clear();
}
// 下游:alerts.addSink(alertSink); // 告警链路独立于业务流
五、选型边界:side output vs filter vs union

| 方案 | 遍历次数 | 数据保留 | 类型 | 适用 |
|---|---|---|---|---|
| filter | N 个条件 N 次 | 不命中的丢 | 与主流同 | 单个布尔过滤 |
| side output | 1 次 | 全保留 | 每 tag 独立 | 多路分发(主力) |
| union | 1 次 | 全合并 | 必须一致 | N 条同类流合并 |
判断顺序:只要一个布尔过滤 → filter;多类数据要分走不同下游 → side output;已经分好的流要合回去 → union。侧输出和 union 甚至可以配合:先 side output 拆开、各自处理后 union 合并回主链路。
六、实战避坑清单
- OutputTag 必须匿名类
{},且 static final——类型信息 + 序列化安全,两条红线一次满足; - 分派记得 return:命中旁路后不 return,记录会同时进主流和旁路(双写);
- 迟到数据要 allowedLateness + sideOutputLateData 配套:只配后者,超时迟到数据会被静默丢弃;
- 侧输出流也是普通流:可以继续窗口/process/再侧输出,别把旁路当"死胡同"只做 Sink;
- 别用侧输出替代 keyBy:多路分发和分组聚合是两个维度,侧输出不改变数据分布;
- tag 的粒度要克制:每个 tag 一条独立子流、下游各占算子,tag 过多会显著增加 DAG 复杂度——同类告警共用一个 tag 带 type 字段,别一个类型一个 tag。
七、总结:我的判断
侧输出流的本质,是 Flink 在"数据流"模型上提供的带类型标签的旁路管道:一次遍历、多路分发、记录不丢、类型自由。相比 filter 的多次遍历和 union 的合并语义,它补上了"拆分"这一环——而且是拆分中最优雅的一环。
三条实操建议:
- 多路分发默认侧输出:实时数仓的 ODS→DWD 分流、异常/迟到/告警旁路,全是它的主场;
- 迟到数据必须成对配置:allowedLateness(窗口重触发)+ sideOutputLateData(旁路兜底),漏一个就是数据静默丢失;
- 关注旁路的下游:侧输出流不是终点,每条旁路都要有明确归宿(补算任务/对账队列/告警中心),否则就是"分出去了但没人接"。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/zmj_0817/article/details/164154249




