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

如何使用Flink流式处理JSON?及字符串流按天聚合计数实现咨询

解决Flink流任务按天聚合统计字符串数量的问题

看起来你正在用Flink处理流数据,需要按天统计字符串总数并输出指定格式的JSON。我来一步步帮你实现这个需求:

步骤1:准备时间戳提取与水位线

首先,你的eventStream是DataStream[String],我们需要先从每个字符串中提取事件时间戳(如果数据源不带时间戳,也可以用处理时间,但事件时间统计更准确)。假设每个输入字符串包含可解析的时间字段(比如格式为yyyy-MM-dd HH:mm:ss),先把字符串转换成带时间戳的结构:

import org.apache.flink.api.common.functions.MapFunction
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor
import org.apache.flink.streaming.api.windowing.time.Time

// 设置Flink时间特性为事件时间
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)

// 假设输入字符串格式为 "2018-03-03 10:00:00,任意内容"
val timestampedStream = eventStream
  .map(new MapFunction[String, (String, Long)] {
    override def map(value: String): (String, Long) = {
      val parts = value.split(",")
      // 解析时间戳为毫秒级
      val timestamp = java.time.LocalDateTime.parse(parts(0), java.time.format.DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))
        .atZone(java.time.ZoneId.systemDefault())
        .toInstant()
        .toEpochMilli()
      (parts(1), timestamp)
    }
  })
  // 设置水位线,允许5秒的乱序数据
  .assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor[(String, Long)](Time.seconds(5)) {
    override def extractTimestamp(element: (String, Long)): Long = element._2
  })

如果不需要事件时间,改用处理时间的话,只需把时间特性改为TimeCharacteristic.ProcessingTime,跳过时间戳提取步骤即可。

步骤2:按天聚合统计数量

接下来要做按自然天对齐的滚动窗口统计(比如北京时间0点到24点),Flink默认窗口按epoch时间对齐,所以需要指定时区偏移量:

import org.apache.flink.streaming.api.scala.function.WindowFunction
import org.apache.flink.streaming.api.windowing.windows.TimeWindow
import org.apache.flink.util.Collector
import java.time.{ZoneId, LocalDate}
import java.time.format.DateTimeFormatter

// 全局统计所有字符串数量,用固定key做keyBy
val dailyAggStream = timestampedStream
  .keyBy(_ => "global")
  // 东八区需偏移-8小时,让窗口对齐到北京时间0点;UTC时区可省略第二个参数
  .window(TumblingEventTimeWindows.of(Time.days(1), Time.hours(-8)))
  .apply(new WindowFunction[(String, Long), (String, Long), String, TimeWindow] {
    override def apply(key: String, window: TimeWindow, input: Iterable[(String, Long)], out: Collector[(String, Long)]): Unit = {
      // 统计窗口内元素总数
      val count = input.size.toLong
      // 将窗口起始时间转换为日期字符串
      val date = LocalDate.ofInstant(java.time.Instant.ofEpochMilli(window.getStart), ZoneId.systemDefault())
        .format(DateTimeFormatter.ISO_LOCAL_DATE)
      out.collect((date, count))
    }
  })

步骤3:转换为指定JSON格式输出

我们需要把每天的聚合结果整理成你要求的JSON结构,同时维护最近N天的结果(对应示例中的days before参数)。这里用Flink的状态管理来保存历史数据,确保故障恢复时不丢失:

import org.apache.flink.api.common.state.ListStateDescriptor
import org.apache.flink.runtime.state.{FunctionInitializationContext, FunctionSnapshotContext}
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction
import org.apache.flink.streaming.api.functions.sink.SinkFunction
import com.fasterxml.jackson.databind.ObjectMapper
import java.util

// 定义POJO类,方便JSON序列化
case class DailyAgg(date: String, sum: Long)
case class AggregationResult(aggregationType: String, `days before`: Int, aggregates: util.List[DailyAgg])

// 自定义Sink,维护最近N天的聚合结果并输出JSON
class DailyAggSink(val daysBefore: Int) extends SinkFunction[(String, Long)] with CheckpointedFunction {
  private var state: org.apache.flink.api.common.state.ListState[DailyAgg] = _
  private val objectMapper = new ObjectMapper()

  override def invoke(value: (String, Long), context: SinkFunction.Context[_]): Unit = {
    val newAgg = DailyAgg(value._1, value._2)
    // 获取当前状态中的历史聚合数据
    val currentAggs = scala.collection.JavaConverters.asScalaIterator(state.get().iterator()).toList
    // 过滤掉超过daysBefore的旧数据,保留最近N天的结果
    val filteredAggs = currentAggs.filter(agg => {
      val aggDate = LocalDate.parse(agg.date)
      val today = LocalDate.now(ZoneId.systemDefault())
      java.time.temporal.ChronoUnit.DAYS.between(aggDate, today) <= daysBefore
    }) :+ newAgg
    // 去重并按日期排序
    val uniqueSortedAggs = filteredAggs
      .groupBy(_.date)
      .map(_._2.head)
      .sortBy(_.date)
      .toList
    // 更新状态
    state.clear()
    uniqueSortedAggs.foreach(state.add)
    // 序列化并输出JSON
    val result = AggregationResult("day", daysBefore, java.util.Arrays.asList(uniqueSortedAggs:_*))
    println(objectMapper.writeValueAsString(result))
    // 也可替换为输出到Kafka、文件等外部存储
  }

  override def snapshotState(context: FunctionSnapshotContext): Unit = {}

  override def initializeState(context: FunctionInitializationContext): Unit = {
    val desc = new ListStateDescriptor[DailyAgg]("daily-aggs", classOf[DailyAgg])
    state = context.getOperatorStateStore.getListState(desc)
  }
}

// 将聚合结果接入自定义Sink,这里指定days before为2
dailyAggStream.addSink(new DailyAggSink(2))

关键注意点

  • 时间对齐:一定要根据业务时区设置窗口偏移量,否则会出现窗口与自然天错位的情况。
  • 状态管理:用ListState维护历史数据,确保流任务故障恢复后数据不丢失。
  • 去重处理:迟到数据可能导致同一窗口多次触发,需要对同一日期的聚合结果去重,保留最新值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:02:09