You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink中如何将流的统计数据定期发送至另一个流

流处理多数据源周期统计实现方案

核心逻辑

不管你当前使用的是Flink、Spark Streaming还是其他流处理框架,实现思路都是统一的:

  • 读取两类数据源时给每条记录打上来源标识,明确区分是Kafka来源还是文件来源
  • 配置时长为5分钟的周期统计窗口,按来源标识分组执行计数
  • 将统计结果序列化后写入目标流即可

以Flink为例的实现示例

1. 定义带来源标识的数据结构

// 所有读取到的记录统一封装为该结构,携带来源标记
public class SourceData {
    // 原始业务数据
    private String rawData;
    // 来源标记,固定取值为 KAFKA / FILE
    private String sourceType;
    // 记录处理时间戳,用于窗口计算
    private Long processTime;
    // 省略构造方法、getter/setter
}

2. 读取数据并打标记后合并流

// 读取Kafka数据并打标记
DataStream<SourceData> kafkaStream = env.addSource(new FlinkKafkaConsumer<>("your-business-topic", new SimpleStringSchema(), kafkaProps))
    .map(raw -> new SourceData(raw, "KAFKA", System.currentTimeMillis()));

// 读取文件流数据并打标记
DataStream<SourceData> fileStream = env.readTextFile("your-file-path")
    .map(raw -> new SourceData(raw, "FILE", System.currentTimeMillis()));

// 合并两个来源的流
DataStream<SourceData> unionStream = kafkaStream.union(fileStream);

3. 5分钟滚动窗口统计各来源记录数

// 按来源分组,5分钟滚动窗口统计数量
DataStream<Tuple2<String, Long>> statResultStream = unionStream
    .keyBy(SourceData::getSourceType)
    // 如果需要按事件时间统计,替换为事件时间滚动窗口+水印配置即可
    .window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
    .aggregate(new CountAggregateFunc(), new WindowResultFunc());

// 自定义聚合函数
public class CountAggregateFunc implements AggregateFunction<SourceData, Long, Long> {
    @Override
    public Long createAccumulator() { return 0L; }
    @Override
    public Long add(SourceData value, Long acc) { return acc + 1; }
    @Override
    public Long getResult(Long acc) { return acc; }
    @Override
    public Long merge(Long acc1, Long acc2) { return acc1 + acc2; }
}

// 自定义窗口输出函数
public class WindowResultFunc implements WindowFunction<Long, Tuple2<String, Long>, String, TimeWindow> {
    @Override
    public void apply(String sourceType, TimeWindow window, Iterable<Long> counts, Collector<Tuple2<String, Long>> out) {
        Long count = counts.iterator().next();
        out.collect(Tuple2.of(sourceType, count));
    }
}

4. 统计结果写入目标流

// 示例:将统计结果写入Kafka统计topic
statResultStream.map(res -> String.format("{\"source\":\"%s\", \"count\":%d, \"windowEnd\":%d}", res.f0, res.f1, window.maxTimestamp()))
    .addSink(new FlinkKafkaProducer<>("your-stat-topic", new SimpleStringSchema(), kafkaProps));

注意事项

  • 若需要精确的事件时间统计,将处理时间窗口替换为事件时间窗口,同时配置对应水印策略即可
  • 开启流处理框架的 checkpoint 机制,避免任务重启后统计状态丢失,保证统计准确性
  • 统计结果的序列化格式可按需调整为JSON、Avro等业务常用格式

内容的提问来源于stack exchange,提问作者Can DALAY

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.04 16:18:04