Flink高级之侧输出流Side Output原理及代码实现:从OutputTag到多流分发

文章来源声明: 原文作者:starzy; 来源站点:掘金; 原文链接:https://juejin.cn/post/7685581483048878132; 本文基于上述来源整理/加工,觅优补充点评,仅供技术学习交流。版权归原作者所有。
觅优短评

文章把侧输出流的原理、代码与选型边界讲得很透,尤其点出 OutputTag 花括号与迟到数据成对配置两个高频坑,适合做实时数仓分流、异常旁路与迟到补算的 Flink 开发者精读。

> **摘要**:一条实时数据流里总混着正常、异常、迟到、需监控四类数据,传统 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——不丢数据、水位线语义一致。

还有一个容易被忽略的点:侧输出流本身也是普通 DataStreamgetSideOutput(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) 的内部实现是 SideOutputDataOutputOutput 接口的实现),它按 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

在这里插入图片描述

方案遍历次数数据保留类型适用
filterN 个条件 N 次不命中的丢与主流同单个布尔过滤
side output1 次全保留每 tag 独立多路分发(主力)
union1 次全合并必须一致N 条同类流合并

判断顺序:只要一个布尔过滤 → filter;多类数据要分走不同下游 → side output;已经分好的流要合回去 → union。侧输出和 union 甚至可以配合:先 side output 拆开、各自处理后 union 合并回主链路。

六、实战避坑清单

  1. OutputTag 必须匿名类 {},且 static final——类型信息 + 序列化安全,两条红线一次满足;
  2. 分派记得 return:命中旁路后不 return,记录会同时进主流和旁路(双写);
  3. 迟到数据要 allowedLateness + sideOutputLateData 配套:只配后者,超时迟到数据会被静默丢弃;
  4. 侧输出流也是普通流:可以继续窗口/process/再侧输出,别把旁路当"死胡同"只做 Sink;
  5. 别用侧输出替代 keyBy:多路分发和分组聚合是两个维度,侧输出不改变数据分布;
  6. tag 的粒度要克制:每个 tag 一条独立子流、下游各占算子,tag 过多会显著增加 DAG 复杂度——同类告警共用一个 tag 带 type 字段,别一个类型一个 tag。

七、总结:我的判断

侧输出流的本质,是 Flink 在"数据流"模型上提供的带类型标签的旁路管道:一次遍历、多路分发、记录不丢、类型自由。相比 filter 的多次遍历和 union 的合并语义,它补上了"拆分"这一环——而且是拆分中最优雅的一环。

三条实操建议:

  1. 多路分发默认侧输出:实时数仓的 ODS→DWD 分流、异常/迟到/告警旁路,全是它的主场;
  2. 迟到数据必须成对配置:allowedLateness(窗口重触发)+ sideOutputLateData(旁路兜底),漏一个就是数据静默丢失;
  3. 关注旁路的下游:侧输出流不是终点,每条旁路都要有明确归宿(补算任务/对账队列/告警中心),否则就是"分出去了但没人接"。