返利app排行榜的实时计算架构:基于Flink的流处理方案

大家好,我是阿可,微赚淘客系统及省赚客APP创始人,是个冬天不穿秋裤,天冷也要风度的程序猿!

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

一、实时排行榜架构设计与技术选型

1.1 架构分层与数据流向

针对返利app排行榜的多维度需求(商品销量榜、用户返利榜、类目热度榜),设计三层流处理架构,数据流向如下:

  1. 数据接入层:通过Kafka接收两类数据源——用户行为数据(商品点击、下单、返利领取,Topic:rebate_user_behavior)、业务基础数据(商品信息、用户信息,Topic:rebate_biz_data);
  2. 计算处理层:Flink作为核心计算引擎,实现数据清洗、维度关联、指标聚合与排行榜排序;
  3. 结果输出层:将实时计算的排行榜数据写入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;
        }
    }
}

三、架构优化与监控保障

  1. 状态后端优化:采用RocksDB作为Flink状态后端,将窗口聚合状态持久化到本地磁盘,避免内存溢出,配置state.backend.rocksdb.localdir: /data/flink/rocksdb
  2. 背压处理:通过Flink Web UI监控背压指标,当Kafka数据接入速率超过计算速率时,动态调整并行度(如从3提升至5),或开启Kafka消费者限流(flink.kafka.consumer.max.poll.records: 1000);
  3. 监控告警:集成Prometheus+Grafana监控Flink Job指标(如窗口延迟、数据处理量),当“窗口延迟超过30秒”时触发企业微信告警,及时排查计算瓶颈;
  4. 数据一致性:对关键排行榜(如商品销量榜),定期通过Flink批处理Job与MySQL离线数据对账,确保实时计算结果误差率≤0.1%。

本文著作权归聚娃科技省赚客app开发者团队,转载请注明出处!

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐