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

Flink用户指纹流处理作业并发问题:重复归因事件的排查与解决方案咨询

解决Flink用户指纹识别作业中未归因事件重复归因的问题

你正在运行一个基于Apache Flink的点击流事件用户指纹识别作业,核心逻辑是对事件进行多维度归因处理,但目前遇到了部分未归因事件被重复处理生成重复归因事件的问题,重复事件的时间戳恰好相差30秒(即你设置的TumblingProcessingTimeWindow窗口大小)。你推测问题出在AttributeBackLogEvents ProcessFunction中,并发任务读取MySQL中的同一批未归因事件导致重复处理,尝试过forceNonParallel()但未解决,考虑过select for update但担心死锁,希望找到可行的解决方案。


核心作业代码

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// setting event time characteristic for processing
env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime);

DataStream<EventData> input = ConfluentKafkaSource.createKafkaSourceFromApplicationProperties(env);

final OutputTag<EventData> emailPresentTag = new OutputTag<>("email-present") { };
final OutputTag<EventData> dispatchIdPresentTag = new OutputTag<>("dispatch-id-present") { };
final OutputTag<EventData> residueTag = new OutputTag<>("residue") { };

SingleOutputStreamOperator<EventData> splitStream = input
    .process(new ProcessFunction<EventData, EventData>() {
        @Override
        public void processElement(EventData data, Context ctx, Collector<EventData> out) {
            if (data.email != null && !data.email.isEmpty()) {
                // emit data to side output for emailPresentTag
                ctx.output(emailPresentTag, data);
            } else if (data.url != null && data.url.contains("utm_source=starling")) {
                // emit data to side output for dispatchIdPresentTag
                ctx.output(dispatchIdPresentTag, data);
            } else {
                // emit data to side output for ip/campaign attributing
                ctx.output(residueTag, data);
            }
        }
    });

DataStream<EventData> emailPresentStream = splitStream.getSideOutput(emailPresentTag);
DataStream<EventData> dispatchIdPresentStream = splitStream.getSideOutput(dispatchIdPresentTag);
DataStream<EventData> residueStream = splitStream.getSideOutput(residueTag);

// process the 3 split streams separately based on their corresponding logic
DataStream<EventData> enrichedEmailPresentStream = emailPresentStream
    .keyBy(e -> e.lbUserId == null ? e.eventId : e.lbUserId)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(30)))
    .process(new AttributeWithEmailPresent());

DataStream<EventData> enrichedDispatchIdPresentStream = dispatchIdPresentStream
    .keyBy(e -> e.lbUserId == null ? e.eventId : e.lbUserId)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(30)))
    .process(new AttributeWithDispatchPresent());

DataStream<EventData> enrichedResidueStream = residueStream
    .keyBy(e -> e.lbUserId == null ? e.eventId : e.lbUserId)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(30)))
    .process(new AttributeWithIP());

DataStream<EventData> dataStream = enrichedEmailPresentStream.union(enrichedDispatchIdPresentStream, enrichedResidueStream);

final OutputTag<EventData> attributedTag = new OutputTag<>("attributed") { };
final OutputTag<EventData> unattributedTag = new OutputTag<>("unattributedTag") { };

SingleOutputStreamOperator<EventData> splitEnrichedStream = dataStream
    .process(new ProcessFunction<EventData, EventData>() {
        @Override
        public void processElement(EventData data, Context ctx, Collector<EventData> out) {
            if (data.attributedEmail != null && !data.attributedEmail.isEmpty()) {
                // emit data to side output for emailPresentTag
                ctx.output(attributedTag, data);
            } else {
                // emit data to side output for ip/campaign attributing
                ctx.output(unattributedTag, data);
            }
        }
    });

//splitting attributed and unattributed stream
DataStream<EventData> attributedStream = splitEnrichedStream.getSideOutput(attributedTag);
DataStream<EventData> unattributedStream = splitEnrichedStream.getSideOutput(unattributedTag);

// attributing backlog unattributed events using attributed stream and flushing resultant attributed
// stream to kafka enriched_clickstream_event topic.
attributedStream = attributedStream
    .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(30)))
    .process(new AttributeBackLogEvents())
    .forceNonParallel();

attributedStream
    .addSink(ConfluentKafkaSink.createKafkaSinkFromApplicationProperties())
    .name("Enriched Event kafka topic sink");

//handling unattributed events. Flushing them to mysql
Properties dbProperties = ConfigReader.getConfig().get(REPORTINGDB_PREFIX);
ObjectMapper objectMapper = new ObjectMapper();
unattributedStream.addSink(JdbcSink.sink(
        "INSERT IGNORE INTO events_store.unattributed_event (event_id, lb_user_id, ip, event) values (?,?,?,?)",
        (ps, t) -> {
            ps.setString(1, t.eventId);
            ps.setString(2, t.lbUserId);
            ps.setString(3, t.ip);
            try {
                ps.setString(4, objectMapper.writeValueAsString(t));
            } catch (JsonProcessingException e) {
                logger.error("[UserFingerPrintJob] "+ e.getMessage());
            }
        },
        JdbcExecutionOptions.builder()
                .withBatchIntervalMs(Long.parseLong(dbProperties.getProperty(REPORTINGDB_FLUSH_INTERVAL)))
                .withMaxRetries(Integer.parseInt(dbProperties.getProperty(REPORTINGDB_FLUSH_MAX_RETRIES)))
                .build(),
        new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
                .withUrl(dbProperties.getProperty(REPORTINGDB_URL_PROPERTY_NAME))
                .withDriverName(dbProperties.getProperty(REPORTINGDB_DRIVER_PROPERTY_NAME))
                .withUsername(dbProperties.getProperty(REPORTINGDB_USER_PROPERTY_NAME))
                .withPassword(dbProperties.getProperty(REPORTINGDB_PASSWORD_PROPERTY_NAME))
                .build())).name("Unattributed event ReportingDB sink");

env.execute("UserFingerPrintJob");

作业执行流程

  • 步骤1:根据事件的email存在性、URL包含特定utm参数、剩余事件三个规则,将输入流拆分为三个子流,分别进行归因处理后合并为一个流
  • 步骤2:将合并后的流拆分为已归因流和未归因流,未归因事件写入MySQL的unattributed_event表作为积压数据
  • 步骤3:已归因事件进入AttributeBackLogEvents ProcessFunction,从MySQL读取与当前事件同lb_user_id(Cookie ID)或同IP的未归因积压事件,完成归因后与当前已归因事件一同输出到Kafka的enriched_clickstream_event主题

问题分析

你观察到重复的归因事件时间戳相差30秒,正好对应窗口的大小,说明同一批未归因事件被多个窗口周期的AttributeBackLogEvents任务读取并处理。虽然你用了forceNonParallel(),但可能存在以下问题:

  • 窗口触发时,前一个窗口的任务还未完成处理并标记事件为已处理,下一个窗口的任务又读取了同一批事件
  • forceNonParallel()可能未完全生效(比如上游算子并行度导致数据重复分发,或者Flink的状态管理问题)

可行解决方案

方案1:MySQL行级锁+状态标记(推荐)

这是最可靠的方案,通过数据库层面的状态标记和行级锁,确保每个事件只会被处理一次,同时避免死锁。

具体步骤:

  1. 修改MySQL表结构:给unattributed_event表添加状态和时间字段
    ALTER TABLE events_store.unattributed_event
    ADD COLUMN processing_status VARCHAR(20) DEFAULT 'UNPROCESSED',
    ADD COLUMN updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP;
    
  2. 调整AttributeBackLogEvents的查询逻辑:
    • 使用SELECT ... FOR UPDATE SKIP LOCKED(MySQL 8.0+支持)查询未处理的事件,这个语句会锁定符合条件的行,并且跳过已经被其他事务锁定的行,避免等待导致的死锁
    • 示例查询语句:
      SELECT * FROM events_store.unattributed_event
      WHERE processing_status = 'UNPROCESSED'
      AND (lb_user_id = ? OR ip = ?)
      FOR UPDATE SKIP LOCKED;
      
    • 查询到事件后,立即将这些事件的processing_status更新为PROCESSING
    • 完成归因处理后,将状态更新为PROCESSED,或者直接删除这些事件(如果不需要保留积压历史)
  3. 处理失败恢复:添加定时任务(MySQL事件或外部定时脚本),将超过一定时间(比如5分钟)的PROCESSING状态事件重置为UNPROCESSED,避免任务失败导致事件永远无法被处理

优点:

  • 彻底避免重复处理,状态标记明确
  • SKIP LOCKED避免了死锁和事务等待,高并发下性能友好
  • 失败恢复机制完善,不会丢失事件

方案2:优化Flink窗口与单实例处理

如果你不想修改数据库结构,可以尝试优化Flink的处理逻辑,确保AttributeBackLogEvents的单实例处理真正生效,并且每个窗口周期只处理一次积压事件。

具体步骤:

  1. 确认forceNonParallel()生效:检查Flink作业的UI,确认AttributeBackLogEvents算子的并行度确实为1。如果没有生效,手动设置并行度:
    attributedStream = attributedStream
        .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(30)))
        .process(new AttributeBackLogEvents())
        .setParallelism(1); // 强制并行度为1
    
  2. 在AttributeBackLogEvents中记录已处理的事件ID:使用Flink的State(比如MapState)记录每个窗口周期内已经处理过的event_id,避免同一窗口或下一个窗口重复处理
    • 处理事件前先检查状态,如果已经处理过则跳过
  3. 调整窗口触发策略:可以考虑使用ProcessingTimeSessionWindows或者调整窗口的延迟时间,确保前一个窗口的处理完全完成后,下一个窗口才触发

优点:

  • 不需要修改数据库结构,纯Flink层面调整
  • 实现相对简单

缺点:

  • 单实例处理可能成为性能瓶颈,当积压事件较多时,处理速度受限
  • Flink状态如果丢失(比如作业重启),可能会导致重复处理

方案3:分布式锁(Redis)

通过Redis的分布式锁,确保同一事件只能被一个任务实例处理。

具体步骤:

  1. 集成Redis客户端:在AttributeBackLogEvents中添加Redis客户端依赖(比如Jedis或Lettuce)
  2. 处理事件前获取锁:对每个未归因事件的event_id,尝试获取Redis锁,锁的过期时间设置为比处理单个事件的最长时间稍长(比如60秒)
    • 示例代码:
      String lockKey = "unattributed_event_lock:" + eventId;
      Boolean lockAcquired = redisTemplate.opsForValue().setIfAbsent(lockKey, "locked", 60, TimeUnit.SECONDS);
      if (lockAcquired == null || !lockAcquired) {
          // 未获取到锁,跳过该事件
          return;
      }
      
  3. 处理完成后释放锁:归因处理完成后,删除Redis中的锁
  4. 处理锁过期:如果任务处理时间超过锁的过期时间,可能会导致重复处理,需要结合事件的状态标记(比如在事件处理完成后标记为已处理,下次即使获取到锁也跳过)

优点:

  • 不需要修改MySQL表结构
  • 分布式锁可以跨多个Flink任务实例生效

缺点:

  • 增加了Redis的依赖,需要维护Redis服务
  • 锁的过期时间难以精准设置,可能导致重复处理或处理延迟

总结

综合来看,**方案1(MySQL行级锁+状态标记)**是最推荐的方案,它兼顾了可靠性和性能,并且能很好地处理任务失败的场景。如果你的MySQL版本低于8.0,不支持SKIP LOCKED,可以退而求其次使用SELECT ... FOR UPDATE,同时优化事务的执行时间,减少死锁的概率(比如缩小事务范围,尽快提交事务)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 14:12:35