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:已归因事件进入
AttributeBackLogEventsProcessFunction,从MySQL读取与当前事件同lb_user_id(Cookie ID)或同IP的未归因积压事件,完成归因后与当前已归因事件一同输出到Kafka的enriched_clickstream_event主题
问题分析
你观察到重复的归因事件时间戳相差30秒,正好对应窗口的大小,说明同一批未归因事件被多个窗口周期的AttributeBackLogEvents任务读取并处理。虽然你用了forceNonParallel(),但可能存在以下问题:
- 窗口触发时,前一个窗口的任务还未完成处理并标记事件为已处理,下一个窗口的任务又读取了同一批事件
forceNonParallel()可能未完全生效(比如上游算子并行度导致数据重复分发,或者Flink的状态管理问题)
可行解决方案
方案1:MySQL行级锁+状态标记(推荐)
这是最可靠的方案,通过数据库层面的状态标记和行级锁,确保每个事件只会被处理一次,同时避免死锁。
具体步骤:
- 修改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; - 调整
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,或者直接删除这些事件(如果不需要保留积压历史)
- 使用
- 处理失败恢复:添加定时任务(MySQL事件或外部定时脚本),将超过一定时间(比如5分钟)的
PROCESSING状态事件重置为UNPROCESSED,避免任务失败导致事件永远无法被处理
优点:
- 彻底避免重复处理,状态标记明确
SKIP LOCKED避免了死锁和事务等待,高并发下性能友好- 失败恢复机制完善,不会丢失事件
方案2:优化Flink窗口与单实例处理
如果你不想修改数据库结构,可以尝试优化Flink的处理逻辑,确保AttributeBackLogEvents的单实例处理真正生效,并且每个窗口周期只处理一次积压事件。
具体步骤:
- 确认
forceNonParallel()生效:检查Flink作业的UI,确认AttributeBackLogEvents算子的并行度确实为1。如果没有生效,手动设置并行度:attributedStream = attributedStream .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(30))) .process(new AttributeBackLogEvents()) .setParallelism(1); // 强制并行度为1 - 在
AttributeBackLogEvents中记录已处理的事件ID:使用Flink的State(比如MapState)记录每个窗口周期内已经处理过的event_id,避免同一窗口或下一个窗口重复处理- 处理事件前先检查状态,如果已经处理过则跳过
- 调整窗口触发策略:可以考虑使用
ProcessingTimeSessionWindows或者调整窗口的延迟时间,确保前一个窗口的处理完全完成后,下一个窗口才触发
优点:
- 不需要修改数据库结构,纯Flink层面调整
- 实现相对简单
缺点:
- 单实例处理可能成为性能瓶颈,当积压事件较多时,处理速度受限
- Flink状态如果丢失(比如作业重启),可能会导致重复处理
方案3:分布式锁(Redis)
通过Redis的分布式锁,确保同一事件只能被一个任务实例处理。
具体步骤:
- 集成Redis客户端:在
AttributeBackLogEvents中添加Redis客户端依赖(比如Jedis或Lettuce) - 处理事件前获取锁:对每个未归因事件的
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; }
- 示例代码:
- 处理完成后释放锁:归因处理完成后,删除Redis中的锁
- 处理锁过期:如果任务处理时间超过锁的过期时间,可能会导致重复处理,需要结合事件的状态标记(比如在事件处理完成后标记为已处理,下次即使获取到锁也跳过)
优点:
- 不需要修改MySQL表结构
- 分布式锁可以跨多个Flink任务实例生效
缺点:
- 增加了Redis的依赖,需要维护Redis服务
- 锁的过期时间难以精准设置,可能导致重复处理或处理延迟
总结
综合来看,**方案1(MySQL行级锁+状态标记)**是最推荐的方案,它兼顾了可靠性和性能,并且能很好地处理任务失败的场景。如果你的MySQL版本低于8.0,不支持SKIP LOCKED,可以退而求其次使用SELECT ... FOR UPDATE,同时优化事务的执行时间,减少死锁的概率(比如缩小事务范围,尽快提交事务)。
内容的提问来源于stack exchange,提问作者Shubham

