ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

基于Flink流处理的动态实时亿级全端用户画像系统实战

基于Flink流处理的动态实时亿级全端用户画像系统实战 简介本资源为基于Flink流处理的动态实时亿级全端用户画像系统完整项目包面向计算机、软件工程、人工智能等专业的在校学生与教师可用于毕业设计、课程设计、项目立项演示或进阶学习。项目围绕实时流计算与用户画像构建展开涵盖数据采集、标签计算与画像存储等核心环节适合具备一定Java与大数据基础的学习者参考。压缩包共327个文件约6.07MB以258个Java源码为主体辅以properties配置、xml与yaml环境文件、sql建表脚本、jar依赖及md说明文档另含词典与停用词资源目录结构清晰便于按模块阅读与二次开发。目前已有309人学习下载。代码经测试可正常运行读者可据此理解Flink实时处理链路、画像标签体系与工程组织方式并在此基础上修改扩展功能用于毕设、课设或作业提交。1. 从一份毕业设计说起Flink 流处理怎么撑起亿级用户画像电商大促零点刚过推荐位要立刻切到「刚加购未付款」的人群风控要同步识别「短时间多设备登录」的账号运营后台要实时看到「近 5 分钟下单用户的地域分布」。这三件事背后是同一个东西用户画像。区别在于传统画像靠 T1 跑批第二天才更新标签而动态实时画像要求标签在秒级内跟着行为变。这份「基于 Flink 流处理的动态实时亿级全端用户画像系统」的毕业设计讲的正是后者——用 Flink 做流处理引擎把 App、小程序、H5、PC 全端埋点汇聚成实时标签再对外提供查询。它适合两类人一是做大数据毕业设计、需要一套能跑通、能讲清架构的学生二是刚接触 Flink 实时计算、想找一个完整场景练手的工程师。源码、数据集、文档三件套的价值不在于「能交差」而在于它把 Kafka、Flink、HBase/Redis、ClickHouse 这条链路串成了一个闭环你能顺着它把「实时标签到底怎么算出来」这件事摸一遍。2. 拆解这套画像系统的技术选型为什么是 Flink 而不是 Spark Streaming2.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分区数不是拍脑袋。按单分区 510 MB/s 的消费能力估算假设 App 端峰值 80 MB/s12 分区留了余量。副本因子至少 2毕业设计环境单机可以设 1但生产必须 2 以上。这里有个坑分区数一旦定了后期扩分区会打乱 key 的顺序性如果下游按 userId 做 keyBy扩分区后同一用户可能落到不同分区状态就散了。所以宁可初期多分几个。2.3 标签计算层Flink 作业的算子链设计数据进了 KafkaFlink 作业要做的事分四步解析、清洗、打标签、写存储。解析层把 JSON 拍平成字段清洗层过滤掉测试账号、机器人流量打标签层是核心按规则或模型给用户打上标签写入层把结果落到 HBase 或 Redis 供查询。打标签有两种模式。规则型标签比如「近 7 天登录次数 3」用 Flink 的 KeyedProcessFunction 加定时器就能算统计型标签比如「近 30 天客单价」需要窗口聚合。毕业设计里通常两种都会涉及这也是它比单纯词频统计复杂的地方。// 规则型标签示例统计用户近 5 分钟内的下单次数超过 2 次打「高频下单」标签 DataStreamUserTag 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); DataStreamString rawStream env.addSource( new FlinkKafkaConsumer(ods_app_event, new SimpleStringSchema(), kafkaProps)); DataStreamUserTag tags rawStream .map(new JsonParser()) // 解析 JSON提取 userId、eventType、timestamp .assignTimestampsAndWatermarks( WatermarkStrategy.EventforBoundedOutOfOrderness(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_serverslocalhost:9092, value_serializerlambda v: json.dumps(v).encode(utf-8)) event_types [click, view, order, login] for i in range(10000): event { userId: fuser_{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, valueevent) 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全量标签、稀疏存储1050 ms中rowkey 设计不当会热点Redis热标签、Top N 用户 1 ms高内存淘汰策略要配好ClickHouse标签圈人、分析100 ms秒级低不适合高并发点查我一般用 HBase 存全量Redis 存最近活跃的几百万用户热标签查询时先查 Redismiss 了再查 HBase 回填。这样兼顾成本和延迟。4.2 HBase rowkey 设计与写入优化HBase 的 rowkey 是灵魂。画像场景常用userId tagId或tagId userId。如果查询模式是「查某用户所有标签」用userId打头如果是「查某标签下所有用户」用tagId打头。毕业设计里两种查询都有可以建两张表用 Flink 双写。// Flink 写 HBase 的 sink 示例 public class HBaseSink extends RichSinkFunctionUserTag { 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(); ValueStateDescriptorUserTag descriptor new ValueStateDescriptor(userTag, UserTag.class); descriptor.enableTimeToLive(ttlConfig);NeverReturnExpired保证过期状态不会被读到OnCreateAndWrite表示每次写入刷新 TTL。注意 TTL 不是实时清理是惰性删除内存不会立刻降但增长会放缓。5.2 现象Kafka 消费延迟越来越高但 CPU 和内存都不高八成是数据倾斜。某个 userId 是测试账号产生了百万级事件全落到一个 keyBy 分区。解决是在 keyBy 前加随机前缀打散或者单独过滤掉异常账号。// 加盐打散userId 前拼 09 随机数聚合后再去掉 DataStreamEvent 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 写入报 RegionTooBusyExceptionrowkey 热点。所有 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 延迟曲线两条线都平稳了才敢接生产流量。希望帮到你。本文还有配套的精品资源点击获取
返回列表