大数据技术实践者头像
关注
Flink高级之侧输出流Side Output原理及代码实现:从OutputTag到多流分发封面图

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

摘要:一条实时数据流里总混着正常、异常、迟到、需监控四类数据,传统 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") 在运行时根本不知道元素类型是什么。解决办法是匿名子类:

// ✅ 匿名子类:子类继承时泛型被固化,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) 的内部实现是 SideOutputDataOutputOutput 接口的实现),它按 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

在这里插入图片描述

方案遍历次数数据保留类型适用
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. 关注旁路的下游:侧输出流不是终点,每条旁路都要有明确归宿(补算任务/对账队列/告警中心),否则就是"分出去了但没人接"。

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/zmj_0817/article/details/164154249

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--