Flink写入MySQL仅保留最新点击触发推荐数据的实现咨询
具体实现方案
核心思路是绕开Flink侧执行删除操作的限制,靠表结构设计+写入逻辑+后置清理满足「仅存用户最新一次触发的推荐结果」的要求,常用落地方式有两种:
- 单用户单次推荐仅返回1个商品的场景:直接把MySQL推荐表的主键设置为
user_id,Flink JDBC Sink开启upsert写入模式(在Flink SQL建表时声明user_id为主键,Connector会自动生成INSERT ... ON DUPLICATE KEY UPDATE语法的写入语句)。相同user_id的新结果写入时,会直接覆盖该用户之前的旧记录,全程不需要执行删除操作,表内天然留存每个用户最新一次触发的推荐结果。 - 单用户单次推荐返回多个带排序的商品(即同一个user_id对应多条item_id、rank记录)的场景:给表新增
trigger_ts字段,存每次点击触发计算时生成的毫秒级时间戳作为本次推荐的版本号,同一次计算生成的所有推荐记录打同一个trigger_ts值。写入时直接批量upsert本次新版本的所有记录即可,不需要同步删旧数据;另外在MySQL侧配置低峰期定时任务,每次执行时按user_id分组,保留每个user_id对应trigger_ts最大的一批记录,删除其余旧版本记录即可,定时任务的执行频率可以按业务对数据一致性的要求设置为1分钟到1小时不等,对业务完全无侵入。
注意不要尝试在Flink侧做删除操作,一来官方JDBC Connector本身对流删除语义的支持有限,二来流计算里同步执行删除会大幅拉高写入延迟,增加链路出错概率,把清理逻辑下沉到存储侧是这类场景的通用做法。
架构选型合理性说明
- Flink实时计算用户行为触发的推荐结果、结果写入MySQL供线上业务查询,是实时推荐场景非常成熟的主流选型。Flink的状态计算、Exactly-Once语义可以保证推荐结果计算的准确性和低延迟,MySQL作为线上服务的结果存储,运维成熟度高、查询延迟稳定,在总数据量千万级以内、查询QPS万级以内的业务场景下,性价比远高于其他复杂架构。
- 如果业务规模进一步增长,比如单用户推荐结果条目多、总数据量过亿、查询QPS超过数万,可以把结果层替换为Redis、HBase这类KV存储,这类存储原生支持单行/批量的版本覆盖,写入和清理效率更高,更适配超大规模的线上推荐场景,但常规业务规模下Flink+MySQL的组合完全适用,不属于选型错误。
内容的提问来源于stack exchange,提问作者chris
相关产品推荐
相关产品推荐

