返利app排行榜的实时计算架构:基于Flink的流处理方案
·
返利app排行榜的实时计算架构:基于Flink的流处理方案
大家好,我是阿可,微赚淘客系统及省赚客APP创始人,是个冬天不穿秋裤,天冷也要风度的程序猿!
在返利app中,商品销量排行榜、用户返利金额排行榜是核心流量入口——用户依赖排行榜快速找到高返利、高销量的商品,平台通过排行榜提升商品转化与用户留存。传统批处理模式(如每日凌晨离线计算)存在24小时延迟,无法满足“大促期间实时更新热门商品”的需求。基于此,我们采用Apache Flink流处理框架,构建“数据接入-实时计算-结果输出”的端到端架构,实现排行榜10秒级更新,核心指标计算延迟控制在500ms以内。以下从架构设计、核心计算逻辑、代码实现三方面展开,附完整技术方案。

一、实时排行榜架构设计与技术选型
1.1 架构分层与数据流向
针对返利app排行榜的多维度需求(商品销量榜、用户返利榜、类目热度榜),设计三层流处理架构,数据流向如下:
- 数据接入层:通过Kafka接收两类数据源——用户行为数据(商品点击、下单、返利领取,Topic:
rebate_user_behavior)、业务基础数据(商品信息、用户信息,Topic:rebate_biz_data); - 计算处理层:Flink作为核心计算引擎,实现数据清洗、维度关联、指标聚合与排行榜排序;
- 结果输出层:将实时计算的排行榜数据写入Redis(供APP前端查询)与Elasticsearch(供后台分析),同时通过Flink CDC同步至MySQL做持久化备份。
1.2 核心技术选型
- 流处理引擎:Apache Flink 1.17(支持事件时间、状态管理,适合实时聚合场景);
- 消息中间件:Kafka 3.0(高吞吐、低延迟,支撑每秒10万+行为数据接入);
- 存储组件:Redis Cluster(排行榜数据缓存,支持ZSet有序集合排序)、Elasticsearch 8.0(日志与明细数据存储);
- 维度数据同步:Flink CDC 2.4(实时同步MySQL中的商品、用户维度数据,避免维度滞后)。
二、核心排行榜计算逻辑与代码实现
2.1 商品销量排行榜(实时下单量统计)
商品销量排行榜以“近1小时下单量”为排序指标,需实时接收用户下单行为,按商品ID聚合下单次数,代码如下:
package cn.juwatech.rebate.flink.job;
import cn.juwatech.rebate.flink.dto.UserBehaviorDTO;
import cn.juwatech.rebate.flink.sink.RedisRankingSink;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.streaming.connectors.kafka.KafkaSource;
import org.apache.flink.streaming.connectors.kafka.KafkaSourceBuilder;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import com.alibaba.fastjson.JSON;
import java.time.Duration;
import java.util.Properties;
/**
* 商品销量实时排行榜Flink Job(近1小时下单量排序)
*/
public class ProductSalesRankingJob {
public static void main(String[] args) throws Exception {
// 1. 初始化Flink执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(3); // 并行度设置(根据Kafka分区数调整)
// 2. 配置Kafka数据源(读取用户行为数据)
String kafkaBootstrapServers = "kafka-node1:9092,kafka-node2:9092";
String topic = "rebate_user_behavior";
String groupId = "product_sales_ranking_group";
KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
.setBootstrapServers(kafkaBootstrapServers)
.setTopics(topic)
.setGroupId(groupId)
.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
.build();
// 3. 读取Kafka数据并转换为DTO
DataStream<String> kafkaStream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source");
SingleOutputStreamOperator<UserBehaviorDTO> behaviorStream = kafkaStream
.map((MapFunction<String, UserBehaviorDTO>) jsonStr -> JSON.parseObject(jsonStr, UserBehaviorDTO.class))
// 过滤下单行为(仅保留behaviorType=PURCHASE的数据)
.filter(behavior -> "PURCHASE".equals(behavior.getBehaviorType()));
// 4. 事件时间水印(解决数据乱序问题,允许5秒延迟)
SingleOutputStreamOperator<UserBehaviorDTO> watermarkStream = behaviorStream
.assignTimestampsAndWatermarks(
WatermarkStrategy.<UserBehaviorDTO>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((behavior, timestamp) -> behavior.getBehaviorTime())
);
// 5. 按商品ID分组,1小时滚动窗口聚合下单量
SingleOutputStreamOperator<Tuple2<String, Long>> salesCountStream = watermarkStream
.keyBy(UserBehaviorDTO::getProductId) // 按商品ID分组
.window(TumblingEventTimeWindows.of(Time.hours(1))) // 1小时滚动窗口
.sum("behaviorCount") // 聚合下单次数(behaviorCount=1表示单次下单)
.map(behavior -> Tuple2.of(behavior.getProductId(), behavior.getBehaviorCount()));
// 6. 将计算结果写入Redis(ZSet存储,实现排行榜排序)
salesCountStream.addSink(new RedisRankingSink(
"redis-node1:6379,redis-node2:6379",
"product_sales_ranking:1h", // Redis Key:商品销量榜(近1小时)
3600 // 过期时间(秒),与窗口时长一致
)).name("Product Sales Redis Sink");
// 7. 执行Flink Job
env.execute("Product Sales Real-Time Ranking Job");
}
}
2.2 Redis排行榜Sink实现(ZSet有序集合)
自定义Flink Sink,将商品销量数据写入Redis ZSet,利用ZSet的score排序特性实现实时排行榜,代码如下:
package cn.juwatech.rebate.flink.sink;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import redis.clients.jedis.HostAndPort;
import redis.clients.jedis.JedisCluster;
import java.util.HashSet;
import java.util.Set;
/**
* 排行榜Redis Sink(基于ZSet存储)
*/
public class RedisRankingSink extends RichSinkFunction<Tuple2<String, Long>> {
private JedisCluster jedisCluster;
private final String redisNodes; // Redis集群节点(格式:host1:port1,host2:port2)
private final String redisKey; // Redis Key(排行榜标识)
private final int expireSeconds; // 过期时间(秒)
public RedisRankingSink(String redisNodes, String redisKey, int expireSeconds) {
this.redisNodes = redisNodes;
this.redisKey = redisKey;
this.expireSeconds = expireSeconds;
}
// 初始化Redis集群连接(生命周期:Job启动时执行一次)
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
Set<HostAndPort> nodeSet = new HashSet<>();
String[] nodes = redisNodes.split(",");
for (String node : nodes) {
String[] hostPort = node.split(":");
nodeSet.add(new HostAndPort(hostPort[0], Integer.parseInt(hostPort[1])));
}
jedisCluster = new JedisCluster(nodeSet);
// 设置Redis密码(若有)
jedisCluster.auth("RebateRedis@2024");
}
// 处理每条计算结果(写入Redis ZSet)
@Override
public void invoke(Tuple2<String, Long> value, Context context) throws Exception {
String productId = value.f0; // 商品ID(ZSet的member)
Long salesCount = value.f1; // 下单量(ZSet的score)
// 1. 写入ZSet:score为下单量,按score降序排序
jedisCluster.zadd(redisKey, salesCount, productId);
// 2. 设置过期时间(避免数据堆积)
jedisCluster.expire(redisKey, expireSeconds);
// 3. 日志打印(可选,用于调试)
System.out.printf("商品ID:%s,近1小时下单量:%d,已写入排行榜%n", productId, salesCount);
}
// 关闭Redis连接(生命周期:Job停止时执行一次)
@Override
public void close() throws Exception {
super.close();
if (jedisCluster != null) {
jedisCluster.close();
}
}
}
2.3 用户返利金额排行榜(多维度聚合)
用户返利排行榜需关联“用户信息”维度数据(如用户昵称、头像),并按“近7天累计返利金额”排序,需通过Flink CDC同步MySQL用户表数据,实现维度关联,代码如下:
package cn.juwatech.rebate.flink.job;
import cn.juwatech.rebate.flink.dto.RebateRecordDTO;
import cn.juwatech.rebate.flink.dto.UserDTO;
import com.ververica.cdc.connectors.mysql.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import com.ververica.cdc.debezium.DebeziumSourceFunction;
import org.apache.flink.api.common.functions.JoinFunction;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.connectors.kafka.KafkaSource;
import static org.apache.flink.streaming.connectors.kafka.KafkaSource.builder;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import com.alibaba.fastjson.JSON;
import java.time.Duration;
/**
* 用户返利金额实时排行榜Flink Job(近7天累计返利)
*/
public class UserRebateRankingJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(2);
// 1. 读取Kafka返利记录数据(Topic:rebate_record)
KafkaSource<String> rebateKafkaSource = builder()
.setBootstrapServers("kafka-node1:9092")
.setTopics("rebate_record")
.setGroupId("user_rebate_ranking_group")
.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer")
.build();
SingleOutputStreamOperator<RebateRecordDTO> rebateStream = env.fromSource(rebateKafkaSource, WatermarkStrategy.noWatermarks(), "Rebate Kafka Source")
.map(jsonStr -> JSON.parseObject(jsonStr, RebateRecordDTO.class))
.assignTimestampsAndWatermarks(
WatermarkStrategy.<RebateRecordDTO>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((rebate, timestamp) -> rebate.getRebateTime())
);
// 2. 通过Flink CDC读取MySQL用户表(实时同步用户维度数据)
DebeziumSourceFunction<UserDTO> mysqlSource = MySqlSource.<UserDTO>builder()
.hostname("mysql-node1")
.port(3306)
.databaseList("rebate_db") // 数据库名
.tableList("rebate_db.user_info") // 表名
.username("flink_cdc")
.password("FlinkCDC@2024")
.deserializer(new UserInfoDeserializer()) // 自定义反序列化器
.startupOptions(StartupOptions.initial()) // 初始全量同步,后续增量同步
.build();
DataStream<UserDTO> userStream = env.addSource(mysqlSource, "MySQL CDC Source");
// 3. 按用户ID分组,7天滑动窗口(每1小时更新一次)聚合返利金额
SingleOutputStreamOperator<Tuple3<String, String, Double>> userRebateSumStream = rebateStream
.keyBy(RebateRecordDTO::getUserId)
.window(SlidingEventTimeWindows.of(Time.days(7), Time.hours(1))) // 7天窗口,1小时滑动一次
.sum("rebateAmount") // 聚合累计返利金额
.map(rebate -> Tuple3.of(rebate.getUserId(), "", rebate.getRebateAmount())); // 占位用户昵称
// 4. 关联用户维度数据(补充用户昵称)
SingleOutputStreamOperator<Tuple3<String, String, Double>> joinedStream = userRebateSumStream
.join(userStream)
.where(Tuple3::f0) // 返利流的用户ID(f0)
.equalTo(UserDTO::getUserId) // 用户流的用户ID
.window(SlidingEventTimeWindows.of(Time.hours(1), Time.hours(1)))
.apply((JoinFunction<Tuple3<String, String, Double>, UserDTO, Tuple3<String, String, Double>>)
(rebateTuple, user) -> Tuple3.of(rebateTuple.f0, user.getNickname(), rebateTuple.f2));
// 5. 写入Redis排行榜(ZSet:score=返利金额,member=用户ID:用户昵称)
joinedStream.addSink(new RedisRankingSink(
"redis-node1:6379,redis-node2:6379",
"user_rebate_ranking:7d",
604800 // 7天过期(60*60*24*7)
)).name("User Rebate Redis Sink");
env.execute("User Rebate Real-Time Ranking Job");
}
// 自定义MySQL用户表反序列化器
public static class UserInfoDeserializer implements com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema<UserDTO> {
@Override
public UserDTO deserialize(com.fasterxml.jackson.databind.JsonNode jsonNode) throws Exception {
UserDTO user = new UserDTO();
user.setUserId(jsonNode.get("after").get("user_id").asText());
user.setNickname(jsonNode.get("after").get("nickname").asText());
user.setAvatarUrl(jsonNode.get("after").get("avatar_url").asText());
return user;
}
}
}
三、架构优化与监控保障
- 状态后端优化:采用RocksDB作为Flink状态后端,将窗口聚合状态持久化到本地磁盘,避免内存溢出,配置
state.backend.rocksdb.localdir: /data/flink/rocksdb; - 背压处理:通过Flink Web UI监控背压指标,当Kafka数据接入速率超过计算速率时,动态调整并行度(如从3提升至5),或开启Kafka消费者限流(
flink.kafka.consumer.max.poll.records: 1000); - 监控告警:集成Prometheus+Grafana监控Flink Job指标(如窗口延迟、数据处理量),当“窗口延迟超过30秒”时触发企业微信告警,及时排查计算瓶颈;
- 数据一致性:对关键排行榜(如商品销量榜),定期通过Flink批处理Job与MySQL离线数据对账,确保实时计算结果误差率≤0.1%。
本文著作权归聚娃科技省赚客app开发者团队,转载请注明出处!
更多推荐



所有评论(0)