☰
Flink用户画像与商品推荐系统实战:工程源码解析与实时链路拆解
2026/9/26 14:45:35 网站建设 项目流程

简介:基于Apache Flink的全端用户画像商品推荐系统项目压缩包,面向大数据方向的计算机专业学生与推荐系统开发者,可用于课程设计、毕业设计或实战练手。系统覆盖用户行为数据采集、Flink实时清洗与聚合、动态用户画像构建以及协同过滤、矩阵分解等推荐算法落地,完整呈现从数据处理到推荐展示的闭环工程思路。压缩包共27个文件,以25个Java源码文件为主体,附2个Maven工程配置XML文件,业务模块如analyservice、user-portrait源码结构清晰,便于按模块阅读与二次开发。包体仅24KB,轻量易部署;已有153人学习浏览,适合想快速上手Flink流处理与实时推荐系统的读者参考借鉴。

1. 把 Flink 用户画像做成能跑的商品推荐:这份 zip 里到底有什么

很多想入门实时推荐的工程师,卡住的地方往往不是算法,而是“一份能落地的代码长什么样”。《基于 Flink 全端用户画像商品推荐系统》这份资源,解压后就是一个完整的 Maven 工程user-portrait-master,里面有pom.xml、analyservice模块以及配套的源码目录,覆盖了从行为数据采集、Flink 实时计算、用户标签构建到商品召回排序的完整链路。它的定位不是教学 PPT,而是一套可以导入 IDEA 直接启动的工程骨架,适合做毕业设计、课程设计,也适合想快速搭一套推荐系统 Demo 的工程师拿来改造成生产项目。我会从工程结构、实时链路的核心算子、画像标签的实现方式、推荐算法的接入点以及部署时最容易踩的坑这几个维度来拆这份资源,让你拿到手之后能顺着代码路径走下去,而不是在 pom 依赖里迷路。

2. 从 pom.xml 看技术选型:Flink 版本、依赖和模块边界

2.1 工程结构里藏着的数据流向

把 zip 解压后,首先看pom.xml和user-portrait-master下的目录布局。这个工程不是一个大而全的单体应用,而是按数据处理的职责拆成了公共父模块和analyservice业务模块。analyservice这个名字已经暗示了它主要负责“分析服务”,也就是把用户行为原始日志转化成结构化标签的核心计算逻辑。常见做法是 Maven 多模块结构,父 pom 统一管理依赖版本,子模块各自维护业务代码,方便后续扩展出recommendservice、dataservice之类的独立服务。

<properties> <flink.version>1.13.2</flink.version> <scala.version>2.12</scala.version> <mysql.version>8.0.23</mysql.version> </properties> <dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-jdbc_${scala.version}</artifactId> <version>${flink.version}</version> </dependency> </dependencies>

这里用 1.13.2 是当时比较稳的版本,如果你本机装的是 Flink 1.17 或 1.18,需要同步升级依赖,否则提交作业到集群时会报序列化或兼容性异常。参数上重点关注flink-connector-kafka_${scala.version}这个写法,Scala 版本是 2.12,意味着你的 Java 工程运行时如果依赖了 Scala 2.13 的 Flink 包,就会出现NoSuchMethodError,这类问题大多出现在你本地 scala-library 版本和 Flink 编译用的版本不一致时。

2.2 实时计算作业的骨架:Source、Transform、Sink

在analyservice模块里,核心作业类一般会按照“Source → Transformation → Sink”三段式组织。这份资源的做法是:从 Kafka 读取用户行为日志(浏览、加购、下单),然后经过 Flink 窗口聚合和状态计算生成用户标签,最后把结果写入 MySQL 或 Redis。下面的代码是一个简化版的可运行骨架,对应资源的UserBehaviorAnalysisJob类:

public class UserBehaviorAnalysisJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 开启检查点,保证故障恢复时数据不丢、不重 env.enableCheckpointing(60 * 1000L, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500L); env.getCheckpointConfig().setCheckpointTimeout(60 * 1000L); Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "localhost:9092"); kafkaProps.setProperty("group.id", "user-portrait-group"); FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>( "user_behavior", new SimpleStringSchema(), kafkaProps ); DataStream<String> rawStream = env.addSource(consumer); DataStream<UserBehavior> behaviorStream = rawStream .map(new MapFunction<String, UserBehavior>() { @Override public UserBehavior map(String line) throws Exception { String[] fields = line.split("\t"); return UserBehavior.of( Long.parseLong(fields[0]), Long.parseLong(fields[1]), Integer.parseInt(fields[2]), fields[3], Long.parseLong(fields[4]) ); } }) .returns(TypeInformation.of(UserBehavior.class)); behaviorStream .keyBy(UserBehavior::getUserId) .process(new UserPortraitProcessFunction()) .addSink(new UserTagJdbcSink()); env.execute("user-portrait-analysis-job"); } }

整体逻辑很好理解:先构造执行环境并开启 checkpoint,然后定义 Kafka consumer 订阅user_behaviorTopic,接着把原始字符串解析成UserBehaviorPOJO,再按用户 ID 分组,进入核心的UserPortraitProcessFunction做状态聚合和标签更新,最后写入 MySQL。生产环境可以把 checkpoint 间隔调到 5 分钟,减少频繁持久化给 HDFS 带来的压力;调试阶段建议间隔设小一点,方便快速看到状态恢复的效果。三个参数里,setMinPauseBetweenCheckpoints(500L)是为了防止两次 checkpoint 之间间隔太短导致数据积压,setCheckpointTimeout则是控制单个 checkpoint 最大的执行时间,超时就标记失败。

3. 用户画像标签计算:结合 ProcessFunction 与状态后端做实时更新

3.1 用户行为标签的建模思路

画像系统最核心的问题不是“用什么算法”,而是“标签怎么定义、怎么更新”。这份资源里把标签分成了三类:基础属性标签(性别、年龄、注册时长)、行为偏好标签(30 天内点击最多的品类、最近一次加购时间)以及实时热度标签(当前会话内浏览过的商品 ID 集合)。行为偏好和实时热度都适合用 Flink 状态来维护,因为它们是典型的“有状态计算”——既依赖当前事件,又依赖历史状态。

代码里的UserPortraitProcessFunction就是干这件事的,它继承了KeyedProcessFunction,按用户 ID 划分 Keyed State,每个用户维护一个 ValueState 或 ListState。下面的代码展示了一个计算“最近一次加购时间”标签的状态实现:

public class UserPortraitProcessFunction extends KeyedProcessFunction<Long, UserBehavior, UserTag> { private ValueState<Long> lastCartTimeState; private ValueState<String> categoryState; private MapState<String, Integer> categoryCountState; @Override public void open(Configuration parameters) throws Exception { ValueStateDescriptor<Long> cartDesc = new ValueStateDescriptor<>("lastCartTime", Long.class); lastCartTimeState = getRuntimeContext().getState(cartDesc); MapStateDescriptor<String, Integer> countDesc = new MapStateDescriptor<>("categoryCount", String.class, Integer.class); categoryCountState = getRuntimeContext().getMapState(countDesc); } @Override public void processElement(UserBehavior value, Context ctx, Collector<UserTag> out) throws Exception { if ("cart".equals(value.getBehavior())) { lastCartTimeState.update(value.getTimestamp()); } categoryCountState.put(value.getCategory(), categoryCountState.contains(value.getCategory()) ? categoryCountState.get(value.getCategory()) + 1 : 1); String topCategory = null; int maxCount = 0; for (Map.Entry<String, Integer> entry : categoryCountState.entries()) { if (entry.getValue() > maxCount) { maxCount = entry.getValue(); topCategory = entry.getKey(); } } UserTag tag = new UserTag(); tag.setUserId(value.getUserId()); tag.setTopCategory(topCategory); tag.setLastCartTime(lastCartTimeState.value()); out.collect(tag); } }

这段逻辑的关键点在于MapState的使用。如果你直接在processElement里用一个本地HashMap来累计品类次数,任务运行一段时间后数据会全部丢失,或者出现重复计数,因为算子重启后会从零开始。而 Flink 的MapState是托管状态,配合 checkpoint 能自动持久化,这也是这份资源和那种“单机统计完塞 Redis”山寨实现最大的区别。注意open()方法里的ValueStateDescriptor必须指定状态名称和类型,名称在整个作业里要唯一,否则不同状态之间会互相覆盖。最后out.collect(tag)输出的是一个UserTag对象,你可以把它再写入 Kafka Topic 或直接 Sink 到 Redis,方便下游推荐服务读取。

3.2 状态过期时间与清理策略

实时画像如果状态不设置 TTL,内存会被用户历史数据撑爆。以“30 天内行为偏好”为例,超过 30 天的行为标签其实已经失去了参考价值,但默认情况下 Flink 会把状态保留到作业重启或手动清理。生产上我一般会给每个ValueStateDescriptor设置 TTL:

StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.days(30)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptor<Long> cartDesc = new ValueStateDescriptor<>("lastCartTime", Long.class); cartDesc.enableTimeTTL(ttlConfig);

这里的OnCreateAndWrite表示创建和写入时都会刷新存活时间,适合“每产生一次加购动作就重新算一次”的场景;NeverReturnExpired则保证 Flink 不会把已过期状态的残留值返回给下游,避免了画像里出现上个月的离谱推荐数据。在调试的时候,你可以临时把 TTL 改成Time.hours(1)来验证清理逻辑是否生效,但上线前一定要改回按业务周期配置的值,别把测试配置带到生产。

3.3 标签计算结果落库的常见方式

UserTagJdbcSink需要注意的点是 JDBC 连接实例化时机。很多初学者把DriverManager.getConnection放在每个invoke()调用里,结果连接数爆炸,MySQL 直接被怼挂。正确做法是在JdbcSinkFunction.open()里初始化连接,复用同一个连接实例,并在close()里释放。这份资源的写法类似下面这样:

public class UserTagJdbcSink extends RichSinkFunction<UserTag> { private Connection conn; private PreparedStatement ps; @Override public void open(Configuration parameters) throws Exception { conn = DriverManager.getConnection("jdbc:mysql://localhost:3306/user_portrait?useSSL=false", "root", "root"); ps = conn.prepareStatement( "INSERT INTO user_tags (user_id, top_category, last_cart_time, update_time) " + "VALUES (?, ?, ?, NOW()) ON DUPLICATE KEY UPDATE top_category=VALUES(top_category), last_cart_time=VALUES(last_cart_time)" ); } @Override public void invoke(UserTag value, Context context) throws Exception { ps.setLong(1, value.getUserId()); ps.setString(2, value.getTopCategory()); ps.setLong(3, value.getLastCartTime() == null ? 0L : value.getLastCartTime()); ps.executeUpdate(); } @Override public void close() throws Exception { if (ps != null) ps.close(); if (conn != null) conn.close(); } }

ON DUPLICATE KEY UPDATE在这里是关键,它保证了同一个用户的新标签会覆盖旧标签,而不是无限插入新行。生产环境如果要提升写入吞吐,可以先用批量提交(每攒 1000 条执行一次executeBatch()),再把 MySQL 的rewriteBatchedStatements=true加上,性能能提升五六倍。

4. 商品推荐算法的工程接入:召回、重排和实时兴趣修正

4.1 离线协同过滤与在线召回结合

纯实时推荐做不到冷启动,因为新用户没有行为历史。这份资源的推荐模块采用了“离线先算相似度,在线实时看行为”的混合策略。离线部分用用户行为日志训练协同过滤模型,生成“物品相似度矩阵”存入 Redis;在线部分则根据用户当前的浏览、加购行为,从 Redis 拉取相似商品作为候选集。这样的好处是实时计算不需要重复做复杂的矩阵运算,只需要在 Flink 里维护用户的实时行为列表即可。

协同过滤的训练代码通常是 Spark 或 Flink 批处理作业完成的,但这份资源为了保持单体工程简单,用了一个离线脚本化的计算方式,核心逻辑如下:

# 离线计算商品相似度,结果写入 Redis import pandas as pd from sklearn.metrics.pairwise import cosine_similarity df = pd.read_csv("user_item_behavior.csv") # 列:user_id, item_id, behavior item_matrix = df.pivot_table(index="user_id", columns="item_id", values="behavior", fill_value=0) item_sim = cosine_similarity(item_matrix.T) item_sim_df = pd.DataFrame(item_sim, index=item_matrix.columns, columns=item_matrix.columns)

这段脚本不是项目的运行时部分,而是作为离线产物的补充说明,帮助理解推荐候选集是怎么来的。和 Flink 实时作业配合的流程是:每天凌晨把相似度矩阵批量计算好写入 Redis,白天 Flink 实时作业只需查询 Redis 获取候选商品,再用规则做重排。

4.2 Redis 存取候选商品列表

实时推荐服务不是直接把所有候选商品塞给用户,而是经过“筛选 — 重排 — 截断”三步。筛选阶段会过滤掉用户已经在购物车里或已下单的商品;重排阶段则根据品类偏好加权,比如用户最近 30 天点击最多的品类是“运动户外”,那么就优先把该品类的商品顶上去;最后只取 Top N 返回给前端。下面的代码展示了 Flink 作业里如何通过 Redis 客户端拉取相似商品:

Jedis jedis = new Jedis("localhost", 6379); String key = "similar:item:" + currentItemId; List<String> similarItemIds = jedis.lrange(key, 0, 20);

lrange使用了 Redis 的 List 结构,在离线阶段用rpush写入排序好的相似商品 ID。注意 Redis 的lrange不像数据库那样有复杂的查询条件,它的排序必须在写入时确定,所以离线计算阶段需要用相似度降序排列后再推入 Redis。在线阶段如果用户的行为发生变化,可以把新的行为 ID 继续追加到另一个 List,比如 “user:realtime:cart”,供下游重排逻辑取用。

4.3 实时兴趣修正:重排权重的动态计算

画像系统更新之后,推荐结果不能等下一次模型训练才调整,否则就失去了“实时推荐”的意义。重排权重的调整可以用 Flink CEP 或简单的KeyedProcessFunction来实现:当用户在短时间内连续点击同一个品类的商品超过阈值,就调高该品类在推荐列表里的权重。下面是一个简化版的重排逻辑:

public class RankAdjustFunction extends KeyedProcessFunction<String, BehaviorEvent, RankScore> { private MapState<String, Integer> categoryClickCount; @Override public void processElement(BehaviorEvent value, Context ctx, Collector<RankScore> out) throws Exception { categoryClickCount.put(value.getCategory(), categoryClickCount.contains(value.getCategory()) ? categoryClickCount.get(value.getCategory()) + 1 : 1); int clickCount = categoryClickCount.get(value.getCategory()); if (clickCount >= 3) { RankScore score = new RankScore(); score.setUserId(value.getUserId()); score.setCategory(value.getCategory()); score.setScore(1.5); // 临时调高该类目权重 out.collect(score); } } }

这里的阈值 3 次和权重 1.5 是业务参数,在大促场景下可以降低到 2 次、权重上调到 2.0,因为用户在大促期间的决策速度更快,需要更及时的兴趣捕捉。另外需要注意,categoryClickCount是MapState,如果没有 TTL,点击量会无限累计到很大的值,所以在 open 时务必给状态设置过期时间,比如 1 天或一个 Session 时长。

5. 部署与运维避坑:Flink 作业和工程源码里的五个高频故障

5.1 zip 包解压后文件校验失败,提示压缩包损坏

现象:从网盘或平台下载的 zip 文件解压到一半提示“CRC 校验失败”或“文件头损坏”,甚至有些压缩软件直接报“不可预料的压缩文件末端”。

原因:这类情况大多是传输过程中文件不完整,或压缩包被第三方存储平台二次处理过,比如伪加密标志被置位。还有一种常见情况是浏览器或下载工具断点续传后有缓存残留。

解决:先用WinRAR的“修复压缩文件”功能试试;如果修复无效,重新下载并对比文件大小是否与页面标注一致,也可以用命令行certutil -hashfile user-portrait-master.zip MD5计算哈希值,和源发布方的哈希比对。

5.2 导入 IDEA 后 Maven 依赖报红,flink-connector-jdbc 找不到

现象:pom.xml里flink-connector-jdbc_2.12右侧出现红色波浪线,下拉依赖列表为空。

原因:Flink 1.13 对应的flink-connector-jdbc在 Maven Central 上没有以flink-connector-jdbc_2.12的坐标发布,实际上它是以flink-connector-jdbc_2.12还是flink-connector-jdbc结尾取决于版本。1.11 之前是后者,1.11 之后改成了带 Scala 版本后缀,但某些小版本之间有例外。

解决:最常见做法是加上<scope>provided</scope>之外的显式版本,或者改用flink-table-planner-blink内部自带的 JDBC 依赖。我一般在本地用 1.13.2 时会直接指定为flink-connector-jdbc_2.12:1.13.2,同时在集群的lib目录确认有驱动。

注意:不要把mysql-connector-java打进作业 JAR 里再通过-j提交,那样会出现类加载冲突,正确做法是放在 Flink 的lib目录下。

5.3 作业运行正常但 MySQL 里没有画像数据

现象:Flink 作业整个看板显示running状态,没有任何异常日志,但 MySQL 的user_tags表一直是空的。

原因:这通常是开启了 checkpoint 但 Sink 没有实现CheckpointedFunction,导致数据一直在内存缓冲区,没有真正提交到数据库。还有可能是 JDBC 驱动自动提交被关闭,事务没有被commit()。

解决:如果RichSinkFunction里用的是conn.setAutoCommit(false),那么每次executeUpdate之后要手动conn.commit();或者干脆去掉setAutoCommit(false),让每条数据自动提交。生产上为了吞吐会保留手动提交,但要在invoke()方法末尾判断累计数量,达到批量阈值才commit()。

5.4 Kafka 消费重复或数据丢失

现象:作业重启后,Redis 里的用户标签大量重复,或者部分用户标签缺失。

原因:Kafka 的 offset 提交和 Flink checkpoint 不同步。如果没开启 checkpoint,Flink 是 At-Most-Once 或 At-Least-Once 语义;如果setRestartStrategy配置为不重启,作业失败后 offset 可能没有回滚到正确位置。

解决:把 checkpoint 的CheckpointingMode设置为EXACTLY_ONCE,并给 Kafka consumer 配置setStartFromLatest()或setStartFromEarliest()之外,还要确保enable.auto.commit=false,因为 Flink 会接管 offset 管理。另外,RestartStrategy要配置为FixedDelayRestartStrategyBuilder或FailureRateRestartStrategy,避免单次异常导致作业永久停机。

5.5 本地启动正常,提交到集群后状态恢复失败

现象:本地跑消费几百条数据没问题,提交到 Standalone 集群或 YARN 上跑了半小时后,突然报State was not found或Recovery process failed。

原因:大多是 checkpoint 目录没有为不同作业分别指定路径,多个作业共用了同一个state.checkpoints.dir,导致状态句柄互相覆盖。还有可能是本机代码里用了本地文件系统路径,但集群是分布式路径。

解决:为每个作业配置独立的state.checkpoints.dir,例如hdfs://nameservice/flink/checkpoints/user-portrait-job,并且进入Flink Web UI的 “Checkpoints” 页面确认Latest Completed Checkpoint正在递增。如果还出现状态不兼容,就看看代码里是否修改了ValueStateDescriptor的名称或类型,一旦修改,之前的状态继续使用会反序列化失败。我的习惯是在开发展位环境直接清掉旧状态目录再做验证,生产则要评估兼容性。

6. 验证推荐效果与画像质量:一份随手可用的测试脚本

拿到这份资源,如果只跑通main方法看到作业启动就收工,那收获不算大。我一般会从三个维度做验证:数据完整性验证、推荐效果验证、性能压测验证。

数据完整性验证主要确认画像标签是否准确覆盖目标用户。你可以写一个简单的 SQL,统计当天有画像更新的用户数占活跃用户数的比例,低于 95% 说明状态 TTL 配置不合理或数据源有缺失,优先检查 Kafka Topic 的消费积压情况:

SELECT COUNT(DISTINCT user_id) AS portrait_user_cnt FROM user_tags WHERE update_time >= CURDATE();

推荐效果验证的经典口径是离线 AUC 或在线 CTR 预估。离线方式是把历史上“用户点击过的商品”作为正样本,系统推荐的候选集作为负样本,用 LR 或简单规则模型计算 AUC。在线方式则是看推荐位上的点击率有没有高于旧策略,一般跑两周 AB 实验。没有实验平台的话,可以用一个最朴素的方式:模拟同一用户在 A/B 两类推荐策略下的点击序列,比较人均点击次数。

性能压测验证时,用 Flink 自带的flink run提交作业后,观察 Web UI 里的Backpressure和Idle指标。如果某个算子出现高 Backpressure,说明下游 Sink 写入成为了瓶颈,常见对策是把UserTagJdbcSink改成批量写入。压测时我还习惯在本地用kafka-console-producer模拟高吞吐行为日志,观察从录入到画像更新的端到端延迟:

kafka-console-producer.sh --broker-list localhost:9092 --topic user_behavior

启动后会进入交互式命令行,直接粘贴一行用户ID 商品ID 品类ID 行为类型 时间戳格式的数据即可。结合 Redis 中该用户的标签更新时间,就能算出端到端延迟大概是多少秒。从那以后我每次验证实时推荐作业,都强制把这条链路走一遍:先看 Kafka 消费是否跟上、再看画像表更新计数、最后用 AB 或模拟点击验证推荐列表是否按预期变化。这套流程虽然简单,但至少能把“作业在跑”和“推荐在生效”这两件事分清楚,希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询