资讯详情

资讯详情

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

摘要一条实时数据流里总混着正常、异常、迟到、需监控四类数据传统 filter 方案要遍历 N 遍、逻辑散落 N 个算子Flink 侧输出流Side Output用一次遍历完成多路分发。这篇文章拆透侧输出从 1N 流模型、OutputTag 的类型信息机制为什么必须带花括号、ctx.output 到旁路记录通道的内部实现到窗口迟到数据的专用 APIallowedLateness sideOutputLateData、双流对账异常旁路等完整代码案例最后给出 side output vs filter vs union 的选型边界。关键词Flink 侧输出流、Side Output、OutputTag、ctx.output、getSideOutput、迟到数据、allowedLateness、sideOutputLateData、多流分发、旁路通道、TypeInformation、代码实现一、一条流里混着四类数据怎么办实时业务里一条输入流几乎总是混着正常业务数据、格式异常/校验失败的数据、迟到数据过了窗口触发时间才到、需要监控的告警数据。经典做法是 N 个 filter 串联DataStreamOrdernormalorders.filter(o-isValid(o));DataStreamOrderinvalidorders.filter(o-!isValid(o));问题很明显每条记录被遍历 N 次N 个条件就是 N 倍处理开销而且 filter 不命中的记录直接丢弃——想保留原始报文排障做不到。侧输出流Side Output解决的就是这件事一条流一次遍历按条件拆成 1 条主流 N 条旁路各走各的下游记录不丢。二、侧输出流 1N 模型核心模型就一句话在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 能提取出 TypeInformationStringOutputTagStringlateTagnewOutputTagString(late){};// ❌ 不带花括号类型擦除后拿到的是 Object序列化器无从构造OutputTagStringlateTagnewOutputTag(late);不带花括号的版本在作业提交时往往不报错运行期侧输出第一次反序列化时才炸——类型对不上错误难排查。所以规范是OutputTag 一律new OutputTagX(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 拆成三路异常流带原始报文、延迟流单独标记全部不丢publicclassOrderSplitterextendsProcessFunctionOrder,Order{privatestaticfinalOutputTagOrderINVALID_TAGnewOutputTagOrder(invalid){};privatestaticfinalOutputTagOrderLATE_TAGnewOutputTagOrder(late){};OverridepublicvoidprocessElement(Ordero,Contextctx,CollectorOrderout){// 1) 校验失败 → 异常旁路保留原始数据供排障if(o.orderIdnull||o.amount0){ctx.output(INVALID_TAG,o);return;// 校验失败的不进主流}// 2) 事件时间明显滞后于当前水位 → 延迟旁路longwmctx.timerService().currentWatermark();if(wmLong.MIN_VALUEo.eventTswm-60_000L){ctx.output(LATE_TAG,o);return;}// 3) 正常订单 → 主流继续 DWD 加工out.collect(o);}}DataStreamOrdermainorders.process(newOrderSplitter());DataStreamOrderinvalidmain.getSideOutput(OrderSplitter.INVALID_TAG);DataStreamOrderlatemain.getSideOutput(OrderSplitter.LATE_TAG);// invalid → 对账/人工处理队列late → 延迟补偿链路main → 正常加工注意return的用法分派是互斥的命中旁路就 return避免一条记录同时进主流和旁路。这是侧输出写法里最常见的逻辑 bug——漏了 return数据就双写了。4.2 窗口迟到数据allowedLateness sideOutputLateData这是侧输出最经典的生产场景。窗口触发计算后allowedLateness 允许延迟内的迟到数据重新触发窗口增量更新而 allowedLateness 之外的迟到数据通过sideOutputLateData进旁路——补算或修正不污染已触发的结果OutputTagOrderlateTagnewOutputTagOrder(window-late){};DataStreamWindowResultresultorders.keyBy(Order::getProductId).window(TumblingEventTimeWindows.of(Time.minutes(5))).allowedLateness(Time.minutes(1))// 窗口触发后再等 1 分钟.sideOutputLateData(lateTag)// 1 分钟后还迟到的 → 旁路.aggregate(newAmountAgg(),newWindowStat());// 迟到旁路单独做补算任务T1 修正 / 告警人工介入DataStreamOrderlateOrdersresult.getSideOutput(lateTag);两个关键参数要配套理解allowedLateness(1min)窗口触发后 1 分钟内到达的迟到数据会重新触发窗口计算结果增量更新sideOutputLateData(tag)超过 allowedLateness 的迟到数据不再触发窗口而是进旁路。线上常见的错误是只配了sideOutputLateData忘配allowedLateness——那样迟到数据既不会触发窗口、也进不了旁路直接被静默丢弃补数链路形同虚设。4.3 双流对账异常旁路订单流 支付流对账匹配失败的记录进旁路未匹配订单、未匹配支付分开主流只保留匹配成功的结果publicclassReconcileextendsCoProcessFunctionOrder,Pay,MatchResult{privatestaticfinalOutputTagOrderUNMATCHED_ORDERnewOutputTagOrder(unmatched-order){};privatestaticfinalOutputTagPayUNMATCHED_PAYnewOutputTagPay(unmatched-pay){};privateValueStateOrderpending;OverridepublicvoidprocessElement1(Ordero,Contextctx,CollectorMatchResultout){pending.update(o);ctx.timerService().registerEventTimeTimer(o.ts10*60*1000L);}OverridepublicvoidprocessElement2(Payp,Contextctx,CollectorMatchResultout){Orderopending.value();if(onull){ctx.output(UNMATCHED_PAY,p);// 只有支付没有订单 → 异常旁路}else{pending.clear();out.collect(newMatchResult(o,p,MATCHED));}}OverridepublicvoidonTimer(longts,OnTimerContextctx,CollectorMatchResultout){Orderopending.value();if(o!null){ctx.output(UNMATCHED_ORDER,o);// 10 分钟没匹配 → 订单异常旁路pending.clear();}}}注意两个侧输出 tag 的类型不同UNMATCHED_ORDER 是 Order、UNMATCHED_PAY 是 Pay——侧输出流的类型自由度在这里直接体现两个异常流可以各自接各自的修复链路。4.4 定时器告警旁路在 KeyedProcessFunction 里定时器到期的告警走侧输出业务流保持干净监控逻辑不会因为告警 Sink 出问题而拖垮主流程privatestaticfinalOutputTagAlertALERT_TAGnewOutputTagAlert(alert){};OverridepublicvoidonTimer(longts,OnTimerContextctx,CollectorOrderout)throwsException{Integerccount.value();if(c!nullcexpected){// 告警 → 侧输出流单独接告警中心ctx.output(ALERT_TAG,newAlert(ctx.getCurrentKey(),ts,count-lag));}count.clear();}// 下游alerts.addSink(alertSink); // 告警链路独立于业务流五、选型边界side output vs filter vs union方案遍历次数数据保留类型适用filterN 个条件 N 次不命中的丢与主流同单个布尔过滤side output1 次全保留每 tag 独立多路分发主力union1 次全合并必须一致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旁路兜底漏一个就是数据静默丢失关注旁路的下游侧输出流不是终点每条旁路都要有明确归宿补算任务/对账队列/告警中心否则就是分出去了但没人接。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →