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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:43:14