如何在Flink中用Java对JSON数据做15秒窗口分组聚合统计?
当然可以用Flink实现这个需求!
我会给你一份完整的Java示例代码,完全匹配你想要的输出格式,咱们一步步拆解实现:
1. 先定义数据实体类
首先我们需要一个POJO类来映射JSON数据,这样后续处理会更清晰(如果不想用Lombok,就自己手动写getter、setter和构造方法):
import lombok.Data; @Data public class UserEvent { private String name; private boolean myparam0; private String myparam1; private String myparam2; private String myparam3; private String myparam4; private String ver; }
2. 修正JSON解析逻辑
你的原始输入里是两个JSON对象连在一起,我假设实际是每行一个合法JSON对象(如果不是,需要先做拆分处理,比如按}{分割后补全括号)。这里用Jackson的ObjectMapper来解析,比JSON.parseFull更稳定:
import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class UserParamStatsJob { // 全局复用ObjectMapper,避免重复创建 private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); // 测试阶段用1并行度,生产环境可根据集群调整 // 读取文件并解析为UserEvent数据流 DataStream<UserEvent> inputStream = env.readTextFile("file:///home/ravisankar/workspace/temporary/input.file") .map(line -> OBJECT_MAPPER.readValue(line, UserEvent.class));
3. 分组+15秒时间窗口统计
接下来按name分组,用滚动时间窗口(每15秒统计一次)来做聚合:
import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; // 按name分组,设置15秒滚动窗口,然后聚合统计 inputStream.keyBy(UserEvent::getName) .window(TumblingProcessingTimeWindows.of(Time.seconds(15))) .aggregate(new StatsAggregateFunc(), new StatsWindowResultFunc()) .print(); // 输出结果,生产可替换为写入文件/Kafka等 env.execute("User Param Statistics Job"); } }
如果需要滑动窗口(比如每5秒统计最近15秒的数据),把窗口换成:
.window(SlidingProcessingTimeWindows.of(Time.seconds(15), Time.seconds(5)))
4. 实现聚合逻辑(AggregateFunction)
这个函数负责在窗口内累积统计数据:
import org.apache.flink.api.common.functions.AggregateFunction; import java.util.HashMap; import java.util.Map; // 累加器类,用来临时存储窗口内的统计数据 class StatsAccumulator { private int emptyMyparam3Count; private Map<String, Integer> myparam4CountMap; public StatsAccumulator(int emptyMyparam3Count, Map<String, Integer> myparam4CountMap) { this.emptyMyparam3Count = emptyMyparam3Count; this.myparam4CountMap = myparam4CountMap; } // getter和setter public int getEmptyMyparam3Count() { return emptyMyparam3Count; } public void setEmptyMyparam3Count(int emptyMyparam3Count) { this.emptyMyparam3Count = emptyMyparam3Count; } public Map<String, Integer> getMyparam4CountMap() { return myparam4CountMap; } } // 聚合结果类,存储最终统计值 class StatsResult { private int myparams3; private Map<String, Integer> myparam4; public StatsResult(int myparams3, Map<String, Integer> myparam4) { this.myparams3 = myparams3; this.myparam4 = myparam4; } // getter public int getMyparams3() { return myparams3; } public Map<String, Integer> getMyparam4() { return myparam4; } } // 自定义聚合函数 public class StatsAggregateFunc implements AggregateFunction<UserEvent, StatsAccumulator, StatsResult> { @Override public StatsAccumulator createAccumulator() { // 初始化累加器:myparam3空值计数为0,myparam4统计Map为空 return new StatsAccumulator(0, new HashMap<>()); } @Override public StatsAccumulator add(UserEvent event, StatsAccumulator accumulator) { // 统计myparam3为空的数量 if (event.getMyparam3() == null || event.getMyparam3().isEmpty()) { accumulator.setEmptyMyparam3Count(accumulator.getEmptyMyparam3Count() + 1); } // 统计每个myparam4值出现的次数 String myparam4Value = event.getMyparam4(); accumulator.getMyparam4CountMap().put( myparam4Value, accumulator.getMyparam4CountMap().getOrDefault(myparam4Value, 0) + 1 ); return accumulator; } @Override public StatsResult getResult(StatsAccumulator accumulator) { // 将累加器数据转换为最终结果 return new StatsResult(accumulator.getEmptyMyparam3Count(), accumulator.getMyparam4CountMap()); } @Override public StatsAccumulator merge(StatsAccumulator a, StatsAccumulator b) { // 合并两个累加器(用于窗口合并场景,比如并行度大于1时) int mergedEmptyCount = a.getEmptyMyparam3Count() + b.getEmptyMyparam3Count(); Map<String, Integer> mergedMap = new HashMap<>(a.getMyparam4CountMap()); b.getMyparam4CountMap().forEach((key, value) -> mergedMap.put(key, mergedMap.getOrDefault(key, 0) + value)); return new StatsAccumulator(mergedEmptyCount, mergedMap); } }
5. 转换为目标JSON格式(WindowFunction)
这个函数负责把分组的用户名和聚合结果转换成你需要的JSON结构:
import org.apache.flink.api.java.tuple.Tuple; import org.apache.flink.streaming.api.functions.windowing.WindowFunction; import org.apache.flink.streaming.api.windowing.windows.TimeWindow; import org.apache.flink.util.Collector; import java.util.HashMap; import java.util.Map; public class StatsWindowResultFunc implements WindowFunction<StatsResult, String, Tuple, TimeWindow> { @Override public void apply(Tuple tuple, TimeWindow window, Iterable<StatsResult> input, Collector<String> out) throws Exception { StatsResult statsResult = input.iterator().next(); String userName = tuple.getField(0).toString(); // 构建你需要的输出结构 Map<String, Object> outputMap = new HashMap<>(); Map<String, Object> userStats = new HashMap<>(); userStats.put("myparams3", statsResult.getMyparams3()); userStats.put("myparam4", statsResult.getMyparam4()); outputMap.put(userName, userStats); // 转换为JSON字符串输出 String jsonOutput = UserParamStatsJob.OBJECT_MAPPER.writeValueAsString(outputMap); out.collect(jsonOutput); } }
依赖说明
代码中用到了Jackson来处理JSON,需要在你的pom.xml中添加依赖:
<dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.15.2</version> <!-- 用最新稳定版即可 --> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.1</version> <!-- 匹配你的Flink版本 --> </dependency>
用这个代码处理你给的示例数据,输出会和你要求的结构完全一致:
{ "user1": { "myparams3": 1, "myparam4": { "http://www.seloger.com/erreur-temporaire/binde": 1, "http://www.seloger.com/erreur-temporaire/husk-pie": 1 } } }
内容的提问来源于stack exchange,提问作者Ravi Shanker Reddy
相关产品推荐
相关产品推荐

