如何使用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
相关产品推荐
相关产品推荐

