关键词: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") 在运行时根本不知道元素类型是什么。解决办法是匿名子类:
<span>// ✅ 匿名子类:子类继承时泛型被固化,TypeExtractor 能提取出 TypeInformation<String></span>
OutputTag<String> lateTag = <span>new</span> <span>OutputTag</span><String>(<span>"late"</span>) {};
<span>// ❌ 不带花括号:类型擦除后拿到的是 Object,序列化器无从构造</span>
OutputTag<String> lateTag = <span>new</span> <span>OutputTag</span><>(<span>"late"</span>);
不带花括号的版本在作业提交时往往不报错,运行期侧输出第一次反序列化时才炸——类型对不上,错误难排查。所以规范是: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 拆成三路,异常流带原始报文、延迟流单独标记,全部不丢:
<span>public</span> <span>class</span> <span>OrderSplitter</span> <span>extends</span> <span>ProcessFunction</span><Order, Order> {
<span>private</span> <span>static</span> <span>final</span> OutputTag<Order> INVALID_TAG = <span>new</span> <span>OutputTag</span><Order>(<span>"invalid"</span>) {};
<span>private</span> <span>static</span> <span>final</span> OutputTag<Order> LATE_TAG = <span>new</span> <span>OutputTag</span><Order>(<span>"late"</span>) {};
<span>@Override</span>
<span>public</span> <span>void</span> <span>processElement</span><span>(Order o, Context ctx, Collector<Order> out)</span> {
<span>// 1) 校验失败 → 异常旁路(保留原始数据供排障)</span>
<span>if</span> (o.orderId == <span>null</span> || o.amount <= <span>0</span>) {
ctx.output(INVALID_TAG, o);
<span>return</span>; <span>// 校验失败的不进主流</span>
}
<span>// 2) 事件时间明显滞后于当前水位 → 延迟旁路</span>
<span>long</span> <span>wm</span> <span>=</span> ctx.timerService().currentWatermark();
<span>if</span> (wm > Long.MIN_VALUE && o.eventTs < wm - <span>60_000L</span>) {
ctx.output(LATE_TAG, o);
<span>return</span>;
}
<span>// 3) 正常订单 → 主流(继续 DWD 加工)</span>
out.collect(o);
}
}
DataStream<Order> main = orders.process(<span>new</span> <span>OrderSplitter</span>());
DataStream<Order> invalid = main.getSideOutput(OrderSplitter.INVALID_TAG);
DataStream<Order> late = main.getSideOutput(OrderSplitter.LATE_TAG);
<span>// invalid → 对账/人工处理队列;late → 延迟补偿链路;main → 正常加工</span>
注意 return 的用法:分派是互斥的,命中旁路就 return,避免一条记录同时进主流和旁路。这是侧输出写法里最常见的逻辑 bug——漏了 return,数据就"双写"了。
4.2 窗口迟到数据:allowedLateness + sideOutputLateData
这是侧输出最经典的生产场景。窗口触发计算后,allowedLateness 允许延迟内的迟到数据重新触发窗口(增量更新),而 allowedLateness 之外的迟到数据,通过 sideOutputLateData 进旁路——补算或修正,不污染已触发的结果:
OutputTag<Order> lateTag = <span>new</span> <span>OutputTag</span><Order>(<span>"window-late"</span>) {};
DataStream<WindowResult> result = orders
.keyBy(Order::getProductId)
.window(TumblingEventTimeWindows.of(Time.minutes(<span>5</span>)))
.allowedLateness(Time.minutes(<span>1</span>)) <span>// 窗口触发后再等 1 分钟</span>
.sideOutputLateData(lateTag) <span>// 1 分钟后还迟到的 → 旁路</span>
.aggregate(<span>new</span> <span>AmountAgg</span>(), <span>new</span> <span>WindowStat</span>());
<span>// 迟到旁路:单独做补算任务(T+1 修正 / 告警人工介入)</span>
DataStream<Order> lateOrders = result.getSideOutput(lateTag);
两个关键参数要配套理解:
allowedLateness(1min):窗口触发后 1 分钟内到达的迟到数据,会重新触发窗口计算(结果增量更新);sideOutputLateData(tag):超过 allowedLateness 的迟到数据不再触发窗口,而是进旁路。
线上常见的错误是只配了 sideOutputLateData 忘配 allowedLateness——那样迟到数据既不会触发窗口、也进不了旁路,直接被静默丢弃,补数链路形同虚设。
4.3 双流对账异常旁路
订单流 + 支付流对账,匹配失败的记录进旁路(未匹配订单、未匹配支付分开),主流只保留匹配成功的结果:
<span>public</span> <span>class</span> <span>Reconcile</span> <span>extends</span> <span>CoProcessFunction</span><Order, Pay, MatchResult> {
<span>private</span> <span>static</span> <span>final</span> OutputTag<Order> UNMATCHED_ORDER = <span>new</span> <span>OutputTag</span><Order>(<span>"unmatched-order"</span>) {};
<span>private</span> <span>static</span> <span>final</span> OutputTag<Pay> UNMATCHED_PAY = <span>new</span> <span>OutputTag</span><Pay>(<span>"unmatched-pay"</span>) {};
<span>private</span> ValueState<Order> pending;
<span>@Override</span>
<span>public</span> <span>void</span> <span>processElement1</span><span>(Order o, Context ctx, Collector<MatchResult> out)</span> {
pending.update(o);
ctx.timerService().registerEventTimeTimer(o.ts + <span>10</span> * <span>60</span> * <span>1000L</span>);
}
<span>@Override</span>
<span>public</span> <span>void</span> <span>processElement2</span><span>(Pay p, Context ctx, Collector<MatchResult> out)</span> {
<span>Order</span> <span>o</span> <span>=</span> pending.value();
<span>if</span> (o == <span>null</span>) {
ctx.output(UNMATCHED_PAY, p); <span>// 只有支付没有订单 → 异常旁路</span>
} <span>else</span> {
pending.clear();
out.collect(<span>new</span> <span>MatchResult</span>(o, p, MATCHED));
}
}
<span>@Override</span>
<span>public</span> <span>void</span> <span>onTimer</span><span>(<span>long</span> ts, OnTimerContext ctx, Collector<MatchResult> out)</span> {
<span>Order</span> <span>o</span> <span>=</span> pending.value();
<span>if</span> (o != <span>null</span>) {
ctx.output(UNMATCHED_ORDER, o); <span>// 10 分钟没匹配 → 订单异常旁路</span>
pending.clear();
}
}
}
注意两个侧输出 tag 的类型不同:UNMATCHED_ORDER 是 Order、UNMATCHED_PAY 是 Pay——侧输出流的类型自由度在这里直接体现,两个异常流可以各自接各自的修复链路。
4.4 定时器告警旁路
在 KeyedProcessFunction 里,定时器到期的告警走侧输出,业务流保持干净(监控逻辑不会因为告警 Sink 出问题而拖垮主流程):
<span>private</span> <span>static</span> <span>final</span> OutputTag<Alert> ALERT_TAG = <span>new</span> <span>OutputTag</span><Alert>(<span>"alert"</span>) {};
<span>@Override</span>
<span>public</span> <span>void</span> <span>onTimer</span><span>(<span>long</span> ts, OnTimerContext ctx, Collector<Order> out)</span> <span>throws</span> Exception {
<span>Integer</span> <span>c</span> <span>=</span> count.value();
<span>if</span> (c != <span>null</span> && c < expected) {
<span>// 告警 → 侧输出流,单独接告警中心</span>
ctx.output(ALERT_TAG, <span>new</span> <span>Alert</span>(ctx.getCurrentKey(), ts, <span>"count-lag"</span>));
}
count.clear();
}
<span>// 下游:alerts.addSink(alertSink); // 告警链路独立于业务流</span>
五、选型边界: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(旁路兜底),漏一个就是数据静默丢失;
- 关注旁路的下游:侧输出流不是终点,每条旁路都要有明确归宿(补算任务/对账队列/告警中心),否则就是"分出去了但没人接"。
文章把侧输出流的原理、代码与选型边界讲得很透,尤其点出 OutputTag 花括号与迟到数据成对配置两个高频坑,适合做实时数仓分流、异常旁路与迟到补算的 Flink 开发者精读。