☰
基于Flink流处理的动态实时亿级全端用户画像系统实战
2026/9/26 18:33:29 网站建设 项目流程

简介:本资源为基于Flink流处理的动态实时亿级全端用户画像系统完整项目包,面向计算机、软件工程、人工智能等专业的在校学生与教师,可用于毕业设计、课程设计、项目立项演示或进阶学习。项目围绕实时流计算与用户画像构建展开,涵盖数据采集、标签计算与画像存储等核心环节,适合具备一定Java与大数据基础的学习者参考。压缩包共327个文件,约6.07MB,以258个Java源码为主体,辅以properties配置、xml与yaml环境文件、sql建表脚本、jar依赖及md说明文档,另含词典与停用词资源,目录结构清晰,便于按模块阅读与二次开发。目前已有309人学习下载。代码经测试可正常运行,读者可据此理解Flink实时处理链路、画像标签体系与工程组织方式,并在此基础上修改扩展功能,用于毕设、课设或作业提交。

1. 从一份毕业设计说起:Flink 流处理怎么撑起亿级用户画像

电商大促零点刚过,推荐位要立刻切到「刚加购未付款」的人群,风控要同步识别「短时间多设备登录」的账号,运营后台要实时看到「近 5 分钟下单用户的地域分布」。这三件事背后是同一个东西:用户画像。区别在于,传统画像靠 T+1 跑批,第二天才更新标签;而动态实时画像要求标签在秒级内跟着行为变。这份「基于 Flink 流处理的动态实时亿级全端用户画像系统」的毕业设计,讲的正是后者——用 Flink 做流处理引擎,把 App、小程序、H5、PC 全端埋点汇聚成实时标签,再对外提供查询。它适合两类人:一是做大数据毕业设计、需要一套能跑通、能讲清架构的学生;二是刚接触 Flink 实时计算、想找一个完整场景练手的工程师。源码、数据集、文档三件套的价值不在于「能交差」,而在于它把 Kafka、Flink、HBase/Redis、ClickHouse 这条链路串成了一个闭环,你能顺着它把「实时标签到底怎么算出来」这件事摸一遍。

2. 拆解这套画像系统的技术选型:为什么是 Flink 而不是 Spark Streaming

2.1 实时画像对计算引擎的三个硬要求

先想清楚画像系统到底要什么。第一是低延迟,用户点了「立即购买」,标签「高购买意向」要在几百毫秒内更新,否则推荐位切过去时人已经走了。第二是状态大,一个亿级用户、每人几百个标签,状态规模轻松到 TB 级,引擎必须能扛住大状态并且支持增量 checkpoint。第三是乱序容忍,全端埋点从不同渠道上报,时间戳参差不齐,晚到几分钟的数据不能直接丢。

Spark Streaming 的微批模型在延迟上天然吃亏,批次间隔再小也有秒级抖动;而 Flink 是真正的逐条流处理,事件驱动,延迟能压到毫秒级。更关键的是 Flink 的状态后端(RocksDB)和 Checkpoint 机制,让大状态下的容错成为可能。这不是说 Spark 不行,批处理场景它依然稳,但「动态实时」四个字把天平压向了 Flink。

2.2 全端数据接入:Kafka 主题怎么划分

全端意味着数据源杂。App 端埋点走移动网关,小程序走微信侧回调,H5 走 Nginx 日志,PC 走服务端 SDK。常见做法是统一打到 Kafka,但主题划分有讲究。我一般按「端类型 + 事件大类」拆,比如ods_app_event、ods_mp_event、ods_web_event,而不是所有端塞一个主题。原因是不同端的数据格式、字段完整度、上报频率差异大,混在一起下游解析要写一堆 if-else,还容易因为某个端的数据倾斜拖垮整个消费组。

# 创建三个端类型的 Kafka 主题,分区数按峰值吞吐估算 # App 端量最大,给 12 分区;小程序和 Web 各 6 分区 kafka-topics.sh --create --bootstrap-server localhost:9092 \ --topic ods_app_event --partitions 12 --replication-factor 2 kafka-topics.sh --create --bootstrap-server localhost:9092 \ --topic ods_mp_event --partitions 6 --replication-factor 2 kafka-topics.sh --create --bootstrap-server localhost:9092 \ --topic ods_web_event --partitions 6 --replication-factor 2

分区数不是拍脑袋。按单分区 5~10 MB/s 的消费能力估算,假设 App 端峰值 80 MB/s,12 分区留了余量。副本因子至少 2,毕业设计环境单机可以设 1,但生产必须 2 以上。这里有个坑:分区数一旦定了,后期扩分区会打乱 key 的顺序性,如果下游按 userId 做 keyBy,扩分区后同一用户可能落到不同分区,状态就散了。所以宁可初期多分几个。

2.3 标签计算层:Flink 作业的算子链设计

数据进了 Kafka,Flink 作业要做的事分四步:解析、清洗、打标签、写存储。解析层把 JSON 拍平成字段;清洗层过滤掉测试账号、机器人流量;打标签层是核心,按规则或模型给用户打上标签;写入层把结果落到 HBase 或 Redis 供查询。

打标签有两种模式。规则型标签,比如「近 7 天登录次数 > 3」,用 Flink 的 KeyedProcessFunction 加定时器就能算;统计型标签,比如「近 30 天客单价」,需要窗口聚合。毕业设计里通常两种都会涉及,这也是它比单纯词频统计复杂的地方。

// 规则型标签示例:统计用户近 5 分钟内的下单次数,超过 2 次打「高频下单」标签 DataStream<UserTag> tagStream = eventStream .keyBy(Event::getUserId) .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1))) .aggregate(new OrderCountAggregator()) .filter(count -> count.getCount() > 2) .map(count -> new UserTag(count.getUserId(), "high_freq_order", System.currentTimeMillis())); // 关键参数说明: // 窗口长度 5 分钟,滑动步长 1 分钟,意味着每 1 分钟输出一次近 5 分钟的结果 // 用 EventTime 而非 ProcessingTime,保证乱序数据也能正确归窗 // 水位线设置通常为最大乱序时间,比如 10 秒

这段代码里,SlidingEventTimeWindows的滑动步长决定了标签更新频率。步长 1 分钟意味着标签最迟 1 分钟后更新,如果你要秒级,就得换成KeyedProcessFunction加状态自己维护。aggregate比reduce更适合画像场景,因为聚合逻辑复杂时增量聚合能省状态。水位线是另一个关键,设太小晚到数据被丢,设太大标签输出延迟高,一般按业务能容忍的乱序程度来,10 秒是个常见起点。

3. 从零跑通最小链路:环境搭建与第一个实时标签

3.1 Flink 本地环境与依赖版本对齐

毕业设计最容易翻车的地方不是代码逻辑,是版本。Flink 1.17 和 1.18 的 API 有差异,Kafka 连接器的版本必须和 Flink 主版本对应,Scala 版本也要一致。我一般用 Flink 1.17.2 + Kafka 连接器 1.17.2 + Scala 2.12 这套组合,稳定且资料多。

# 下载并解压 Flink 1.17.2 wget https://archive.apache.org/dist/flink/flink-1.17.2/flink-1.17.2-bin-scala_2.12.tgz tar -xzf flink-1.17.2-bin-scala_2.12.tgz cd flink-1.17.2 # 把 Kafka 连接器 jar 放进 lib 目录 # flink-sql-connector-kafka-1.17.2.jar 和 flink-connector-kafka-1.17.2.jar 都要放 cp ~/downloads/flink-sql-connector-kafka-1.17.2.jar lib/ cp ~/downloads/flink-connector-kafka-1.17.2.jar lib/ # 启动本地集群 ./bin/start-cluster.sh # 验证 Web UI,默认 8081 端口 curl http://localhost:8081

注意flink-sql-connector-kafka和flink-connector-kafka是两个不同的包,前者用于 SQL 作业,后者用于 DataStream API。很多人只放一个,结果要么 SQL 跑不了,要么 DataStream 报 ClassNotFound。另外 Flink 的lib目录和用户代码的pom.xml依赖要一致,别一个用 1.17 一个用 1.18,否则运行时报序列化异常,这种玄学问题排查起来很费时间。

3.2 用 DataStream API 写第一个标签作业

最小可跑的作业:从 Kafka 读 App 埋点,解析出 userId 和 eventType,统计每个用户近 1 分钟的点击次数,超过 5 次输出「活跃用户」标签。

public class ActiveUserTagJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 开启 checkpoint,间隔 10 秒,保证故障恢复 env.enableCheckpointing(10000); // 设置事件时间语义 env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "localhost:9092"); kafkaProps.setProperty("group.id", "active-user-tag"); DataStream<String> rawStream = env.addSource( new FlinkKafkaConsumer<>("ods_app_event", new SimpleStringSchema(), kafkaProps)); DataStream<UserTag> tags = rawStream .map(new JsonParser()) // 解析 JSON,提取 userId、eventType、timestamp .assignTimestampsAndWatermarks( WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, ts) -> event.getTimestamp())) .filter(e -> "click".equals(e.getEventType())) .keyBy(Event::getUserId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new ClickCountAggregator()) .filter(count -> count.getCount() > 5) .map(count -> new UserTag(count.getUserId(), "active_user", System.currentTimeMillis())); // 输出到控制台,实际项目写 HBase 或 Redis tags.print(); env.execute("ActiveUserTagJob"); } }

enableCheckpointing(10000)是保命设置,没有它作业挂了状态全丢。forBoundedOutOfOrderness(Duration.ofSeconds(10))表示容忍 10 秒乱序,超过 10 秒的数据会被丢弃或进侧输出流。TumblingEventTimeWindows是滚动窗口,不重叠,适合统计固定时间段的行为。aggregate里的ClickCountAggregator需要实现AggregateFunction接口,累加器就是一个 Long 计数器。

跑起来后,往 Kafka 发几条测试数据,看控制台有没有输出。如果没输出,先查水位线有没有推进——窗口不触发最常见的原因就是水位线没到窗口结束时间。可以临时把窗口改成 10 秒,水位线改成 1 秒,快速验证逻辑通不通。

3.3 数据集怎么造:模拟全端埋点写入 Kafka

毕业设计给的数据集通常是离线文件,要转成实时流得自己写个生产者。我一般用 Python 脚本模拟,按用户 ID 随机生成点击、浏览、下单事件,带上时间戳发到对应主题。

import json, random, time from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8')) event_types = ['click', 'view', 'order', 'login'] for i in range(10000): event = { 'userId': f'user_{random.randint(1, 1000)}', 'eventType': random.choice(event_types), 'timestamp': int(time.time() * 1000), 'platform': random.choice(['app', 'mp', 'web']) } producer.send('ods_app_event', value=event) time.sleep(0.01) # 控制发送速率,避免打爆本地 Kafka producer.flush()

userId范围控制在 1000 以内,方便观察同一用户的多次事件是否被正确聚合。time.sleep(0.01)是限速,本地环境不限速容易把 Kafka 写满导致消费延迟。时间戳用毫秒,和 Flink 的TimestampAssigner对齐。如果数据集里时间戳是秒级,记得乘 1000,否则水位线永远推不动。

4. 亿级状态下的性能与存储:HBase、Redis 怎么选怎么配

4.1 标签存储的三种方案对比

标签算出来要存,存哪直接影响查询延迟和成本。常见三种:HBase、Redis、ClickHouse。HBase 适合海量稀疏标签,按 rowkey 查单用户快,但范围查询弱;Redis 适合热标签,延迟亚毫秒,但内存成本高,亿级用户全量放 Redis 不现实;ClickHouse 适合标签分析和圈人,但单点查询不如前两者。

存储适用场景单用户查询延迟亿级成本主要坑
HBase全量标签、稀疏存储10~50 ms中rowkey 设计不当会热点
Redis热标签、Top N 用户< 1 ms高内存淘汰策略要配好
ClickHouse标签圈人、分析100 ms~秒级低不适合高并发点查

我一般用 HBase 存全量,Redis 存最近活跃的几百万用户热标签,查询时先查 Redis,miss 了再查 HBase 回填。这样兼顾成本和延迟。

4.2 HBase rowkey 设计与写入优化

HBase 的 rowkey 是灵魂。画像场景常用userId + tagId或tagId + userId。如果查询模式是「查某用户所有标签」,用userId打头;如果是「查某标签下所有用户」,用tagId打头。毕业设计里两种查询都有,可以建两张表,用 Flink 双写。

// Flink 写 HBase 的 sink 示例 public class HBaseSink extends RichSinkFunction<UserTag> { private Connection connection; private BufferedMutator mutator; @Override public void open(Configuration parameters) throws Exception { Configuration conf = HBaseConfiguration.create(); conf.set("hbase.zookeeper.quorum", "localhost"); connection = ConnectionFactory.createConnection(conf); BufferedMutatorParams params = new BufferedMutatorParams(TableName.valueOf("user_tag")); params.writeBufferSize(4 * 1024 * 1024); // 4MB 缓冲 mutator = connection.getBufferedMutator(params); } @Override public void invoke(UserTag tag, Context context) throws Exception { Put put = new Put(Bytes.toBytes(tag.getUserId())); put.addColumn(Bytes.toBytes("info"), Bytes.toBytes(tag.getTagId()), Bytes.toBytes(tag.getTimestamp())); mutator.mutate(put); } @Override public void close() throws Exception { mutator.close(); connection.close(); } }

BufferedMutator比单条Table.put快一个数量级,writeBufferSize设 4MB 是平衡吞吐和延迟的经验值。注意close()里要先关 mutator 再关 connection,顺序反了会丢缓冲数据。另外 HBase 的hbase-site.xml里hbase.hregion.max.filesize别设太小,否则 region 分裂频繁,写入抖动。

4.3 大状态下的 Checkpoint 调优

亿级用户的状态,Checkpoint 是瓶颈。默认的HashMapStateBackend把状态放内存,大状态直接 OOM。必须换EmbeddedRocksDBStateBackend,并且开启增量 Checkpoint。

// 在 Flink 配置文件中设置,或代码里指定 env.setStateBackend(new EmbeddedRocksDBStateBackend(true)); // true 表示增量 checkpoint env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints"); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); // 两次 checkpoint 间隔至少 5 秒 env.getCheckpointConfig().setCheckpointTimeout(600000); // 超时 10 分钟 env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 不并发,避免抢资源

setMinPauseBetweenCheckpoints很关键,不设的话 Checkpoint 一个接一个,正常数据处理被拖慢。setMaxConcurrentCheckpoints(1)也是同理,并发 Checkpoint 在 RocksDB 下容易把 IO 打满。如果 Checkpoint 持续超时,先看 HDFS 写入带宽,再看 RocksDB 的writeBufferSize是不是太小导致频繁 flush。

5. 避坑与排查:这套链路最容易翻车的五个地方

5.1 现象:作业跑几分钟就 OOM,日志显示 RocksDB 内存超限

原因通常是状态没设 TTL,用户标签越积越多,RocksDB 的 block cache 和 write buffer 把内存吃光。解决是给状态加 TTL,比如标签只保留 30 天。

StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.days(30)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptor<UserTag> descriptor = new ValueStateDescriptor<>("userTag", UserTag.class); descriptor.enableTimeToLive(ttlConfig);

NeverReturnExpired保证过期状态不会被读到,OnCreateAndWrite表示每次写入刷新 TTL。注意 TTL 不是实时清理,是惰性删除,内存不会立刻降,但增长会放缓。

5.2 现象:Kafka 消费延迟越来越高,但 CPU 和内存都不高

八成是数据倾斜。某个 userId 是测试账号,产生了百万级事件,全落到一个 keyBy 分区。解决是在 keyBy 前加随机前缀打散,或者单独过滤掉异常账号。

// 加盐打散:userId 前拼 0~9 随机数,聚合后再去掉 DataStream<Event> salted = stream.map(e -> { int salt = new Random().nextInt(10); e.setSaltedKey(salt + "_" + e.getUserId()); return e; }); // 聚合时用 saltedKey,输出前再还原 userId

加盐会多一轮网络传输,但能解决倾斜。如果倾斜来自少数异常账号,直接过滤更划算。

5.3 现象:窗口不触发,数据一直不输出

先查水位线。如果数据源的时间戳是秒级而代码按毫秒解析,水位线永远停在 1970 年。再查并行度,如果 Kafka 分区数是 12 而 Flink 并行度是 1,只有 1 个分区被消费,水位线推进慢。最后查allowedLateness,默认是 0,晚到数据直接丢,窗口可能因为等不到水位线而不触发。

5.4 现象:HBase 写入报 RegionTooBusyException

rowkey 热点。所有 userId 按字典序集中到少数 region。解决是 rowkey 加哈希前缀,比如md5(userId).substring(0,4) + userId,把写入打散到所有 region。代价是范围查询变慢,因为同一用户的标签可能不在同一 region,但画像场景点查为主,可以接受。

5.5 现象:Checkpoint 失败,报 "Could not complete snapshot"

常见原因是 HDFS 权限或空间不足,或者 RocksDB 的本地目录磁盘满。先看 JobManager 日志里的具体异常,如果是AccessControlException就改 HDFS 目录权限;如果是No space left on device就清理 RocksDB 的tmp目录。另外state.backend.rocksdb.localdir别设在/tmp,系统清理会误删。

6. 进阶技巧:用 Flink SQL 做标签规则热更新

DataStream API 写标签逻辑,改一条规则就要重新打包上线,这在画像场景很痛苦。运营今天要「近 3 天登录 2 次」,明天要「近 7 天登录 5 次」,你不可能天天发版。我后来改用 Flink SQL 加规则表的方式:标签规则存在 MySQL 里,Flink SQL 作业定期加载规则,用LOOKUP JOIN关联事件流和规则表,规则变了不用重启作业。

-- 规则表,存在 MySQL,运营后台可改 CREATE TABLE tag_rule ( rule_id STRING, tag_name STRING, event_type STRING, threshold INT, window_minutes INT, PRIMARY KEY (rule_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://localhost:3306/portrait', 'table-name' = 'tag_rule', 'lookup.cache.max-rows' = '1000', 'lookup.cache.ttl' = '60s' ); -- 事件流 CREATE TABLE user_event ( user_id STRING, event_type STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL '10' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'ods_app_event', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ); -- 关联规则并聚合 INSERT INTO user_tag SELECT e.user_id, r.tag_name, COUNT(*) AS cnt FROM user_event e JOIN tag_rule FOR SYSTEM_TIME AS OF e.ts AS r ON e.event_type = r.event_type GROUP BY e.user_id, r.tag_name, TUMBLE(e.ts, INTERVAL '1' MINUTE) HAVING COUNT(*) > MAX(r.threshold);

lookup.cache.ttl设 60 秒,意味着规则变更最多 1 分钟后生效,不用重启作业。FOR SYSTEM_TIME AS OF是时态表关联,保证用事件发生时的规则版本。这个方案的限制是规则不能太复杂,涉及多事件序列的规则还是得回退到 DataStream。但覆盖 80% 的阈值型标签足够了。

验证规则是否生效,最简单的办法是改 MySQL 里的 threshold,然后往 Kafka 发对应事件,看user_tag表有没有新标签。如果没生效,先查lookup.cache.ttl是不是还没过期,再查 JOIN 条件里的event_type是否匹配。我踩过的坑是 MySQL 驱动版本和 Flink JDBC 连接器不兼容,报No suitable driver,换mysql-connector-java-8.0.28就好了。

这套东西做下来,最大的体会是:实时画像的难点不在 Flink 算子写得多花哨,而在状态怎么管、存储怎么配、规则怎么热更新。我现在的习惯是,任何标签上线前先跑 24 小时压测,看 Checkpoint 时长和 Kafka 延迟曲线,两条线都平稳了才敢接生产流量。希望帮到你。

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

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

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

立即咨询