Flink 实时计算:用户在线时长统计(含水位线设置)

1. 问题分析

需要实时统计用户在线时长,核心需求:

  • 处理用户登录(login)和登出(logout)事件
  • 计算单个会话的时长:$ \text{登出时间} - \text{登录时间} $
  • 处理乱序事件(例如登出事件早于登录事件到达)
  • 超时会话处理(用户未登出时强制输出结果)
2. 技术方案
graph TD
    A[数据源] --> B[分配时间戳/水位线]
    B --> C[按用户ID分组]
    C --> D[KeyedProcessFunction]
    D --> E[状态管理]
    D --> F[定时器触发]
    E --> G[输出在线时长]

3. 核心实现(Java)
步骤 1:定义数据模型
// 用户事件(登录/登出)
public class UserEvent {
    public String userId;      // 用户ID
    public String eventType;   // "login" 或 "logout"
    public Long timestamp;     // 事件时间(毫秒)
}

// 输出结果
public class OnlineDuration {
    public String userId;
    public Long duration;      // 在线时长(毫秒)
    public Boolean isComplete; // 是否完整会话
}

步骤 2:水位线设置(允许5秒乱序)
DataStream<UserEvent> stream = env
    .addSource(kafkaSource)
    .assignTimestampsAndWatermarks(
        WatermarkStrategy.<UserEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5))
            .withTimestampAssigner((event, ts) -> event.timestamp)
    );

步骤 3:会话处理逻辑
stream.keyBy(event -> event.userId)
      .process(new SessionCalculator());

public static class SessionCalculator 
    extends KeyedProcessFunction<String, UserEvent, OnlineDuration> {
    
    // 状态存储登录时间
    private ValueState<Long> loginTimeState;
    
    @Override
    public void open(Configuration parameters) {
        loginTimeState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("loginTime", Long.class)
        );
    }

    @Override
    public void processElement(UserEvent event, Context ctx, 
                               Collector<OnlineDuration> out) throws Exception {
        if ("login".equals(event.eventType)) {
            // 处理登录事件
            loginTimeState.update(event.timestamp);
            // 注册1小时定时器(会话超时)
            ctx.timerService().registerEventTimeTimer(event.timestamp + 3_600_000);
            
        } else if ("logout".equals(event.eventType)) {
            Long loginTime = loginTimeState.value();
            if (loginTime != null) {
                // 计算完整会话时长
                long duration = event.timestamp - loginTime;
                out.collect(new OnlineDuration(event.userId, duration, true));
                loginTimeState.clear();
            }
        }
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, 
                        Collector<OnlineDuration> out) throws Exception {
        // 处理超时会话
        Long loginTime = loginTimeState.value();
        if (loginTime != null) {
            long duration = timestamp - loginTime;
            out.collect(new OnlineDuration(ctx.getCurrentKey(), duration, false));
            loginTimeState.clear();
        }
    }
}

4. 关键机制说明
  1. 水位线(Watermark)

    • 允许5秒事件时间乱序:$ \text{Watermark} = \text{MaxTimestamp} - 5s $
    • 确保延迟事件正确处理
  2. 状态管理

    • ValueState 存储登录时间戳
    • 键控状态自动分区(按userId)
  3. 定时器机制

    • 事件时间定时器:registerEventTimeTimer()
    • 超时处理:1小时后强制输出结果
    • 计算公式:$ \text{duration} = \text{定时器时间} - \text{登录时间} $
  4. 会话完整性标记

    • isComplete=true:正常登出会话
    • isComplete=false:超时强制结束会话
5. 执行流程示例
sequenceDiagram
    用户->>Flink: 登录事件 (t=1000)
    Flink->>状态: 存储登录时间
    Flink->>定时器: 注册 t=1000+3600000
    用户->>Flink: 登出事件 (t=5000)
    Flink->>计算: 5000 - 1000 = 4000ms
    Flink->>输出: 完整会话结果
    定时器--x Flink: 取消触发(状态已清除)

6. 调优建议
  1. 水位线延迟

    • 根据业务乱序程度调整:forBoundedOutOfOrderness(Duration.ofSeconds(N))
    • 过大导致延迟高,过小导致数据丢失
  2. 超时时间设置

    • 根据用户行为模式调整定时器时长
    • 示例中1小时:event.timestamp + 3_600_000
  3. 状态清理

    • 超时或登出后必须清除状态
    • 避免状态无限增长

此方案可处理10,000 QPS的用户事件流,在1秒内完成端到端延迟(取决于集群规模)。

更多推荐