天猫返利app的大数据处理架构:用户画像与个性化推荐系统设计

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

在天猫返利app业务中,用户画像与个性化推荐是核心竞争力——通过分析用户浏览、下单、返利领取等行为数据,构建精准用户画像,再基于画像推送高匹配度的天猫商品及返利活动,可将转化率提升30%以上。为此,我们设计了“实时+离线”双轨大数据处理架构,基于Hadoop生态与实时计算引擎,实现从数据采集到推荐落地的全链路闭环。以下从架构分层、用户画像构建、个性化推荐实现三方面展开,附关键代码示例。
在这里插入图片描述

一、大数据处理架构整体设计

架构采用五层分布式设计,各层职责与技术选型如下:

  1. 数据采集层:通过Flume采集用户行为日志(如商品点击、加入购物车),使用Canal同步天猫开放平台API数据(商品价格、返利比例),数据统一写入Kafka消息队列;
  2. 数据存储层:HDFS存储离线原始数据,HBase存储用户画像标签(支持高并发读写),Redis缓存实时推荐结果;
  3. 计算层:离线计算用Spark SQL处理用户行为数据生成画像标签,实时计算用Flink处理用户实时行为(如当前浏览商品)触发动态推荐;
  4. 画像层:构建“用户基础属性+行为偏好+返利敏感度”三维标签体系,通过标签权重算法生成用户画像;
  5. 推荐层:基于协同过滤与标签匹配算法,生成个性化商品推荐列表,通过API接口推送给app前端。

二、用户画像核心组件实现

2.1 数据采集与预处理代码(Flink实时处理)

通过Flink消费Kafka中的用户行为日志,过滤无效数据并提取关键字段,代码如下:

package cn.juwatech.tmallrebate.data.collect;

import cn.juwatech.tmallrebate.dto.UserBehaviorDTO;
import org.apache.flink.api.common.functions.FilterFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import com.alibaba.fastjson.JSON;
import java.util.Properties;

/**
 * 用户行为数据实时采集与预处理
 */
public class UserBehaviorCollector {
    public static void main(String[] args) throws Exception {
        // 1. 初始化Flink执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(3);

        // 2. 配置Kafka消费者(主题:tmall_user_behavior)
        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("bootstrap.servers", "kafka-node1:9092,kafka-node2:9092");
        kafkaProps.setProperty("group.id", "user_behavior_group");
        FlinkKafkaConsumer<String> kafkaConsumer = new FlinkKafkaConsumer<>(
                "tmall_user_behavior",
                new SimpleStringSchema(),
                kafkaProps
        );

        // 3. 读取Kafka数据并转换为DTO
        DataStream<UserBehaviorDTO> behaviorStream = env.addSource(kafkaConsumer)
                .map(jsonStr -> JSON.parseObject(jsonStr, UserBehaviorDTO.class))
                // 4. 过滤无效数据(用户ID为空、行为时间异常)
                .filter((FilterFunction<UserBehaviorDTO>) behavior -> {
                    return behavior.getUserId() != null 
                            && !behavior.getUserId().isEmpty()
                            && behavior.getBehaviorTime() < System.currentTimeMillis();
                });

        // 5. 将预处理后的数据写入HBase(后续用于画像计算)
        behaviorStream.addSink(new HBaseBehaviorSink());

        env.execute("User Behavior Data Collection");
    }
}

// HBase数据写入Sink
class HBaseBehaviorSink extends org.apache.flink.streaming.api.functions.sink.RichSinkFunction<UserBehaviorDTO> {
    private org.apache.hadoop.hbase.client.Connection hbaseConn;
    private org.apache.hadoop.hbase.client.Table behaviorTable;

    @Override
    public void open(org.apache.flink.configuration.Configuration parameters) throws Exception {
        // 初始化HBase连接
        org.apache.hadoop.hbase.conf.Configuration hbaseConf = org.apache.hadoop.hbase.HBaseConfiguration.create();
        hbaseConf.set("hbase.zookeeper.quorum", "zk-node1,zk-node2");
        hbaseConn = org.apache.hadoop.hbase.client.ConnectionFactory.createConnection(hbaseConf);
        behaviorTable = hbaseConn.getTable(org.apache.hadoop.hbase.TableName.valueOf("tmall_user_behavior"));
    }

    @Override
    public void invoke(UserBehaviorDTO behavior, Context context) throws Exception {
        // 构建HBase行键(userId + 行为时间戳)
        String rowKey = behavior.getUserId() + "_" + behavior.getBehaviorTime();
        org.apache.hadoop.hbase.client.Put put = new org.apache.hadoop.hbase.client.Put(org.apache.hadoop.hbase.util.Bytes.toBytes(rowKey));
        // 添加列数据(行为类型、商品ID、返利比例)
        put.addColumn(
                org.apache.hadoop.hbase.util.Bytes.toBytes("cf1"),
                org.apache.hadoop.hbase.util.Bytes.toBytes("behaviorType"),
                org.apache.hadoop.hbase.util.Bytes.toBytes(behavior.getBehaviorType())
        );
        put.addColumn(
                org.apache.hadoop.hbase.util.Bytes.toBytes("cf1"),
                org.apache.hadoop.hbase.util.Bytes.toBytes("productId"),
                org.apache.hadoop.hbase.util.Bytes.toBytes(behavior.getProductId())
        );
        put.addColumn(
                org.apache.hadoop.hbase.util.Bytes.toBytes("cf1"),
                org.apache.hadoop.hbase.util.Bytes.toBytes("rebateRate"),
                org.apache.hadoop.hbase.util.Bytes.toBytes(behavior.getRebateRate().toString())
        );
        behaviorTable.put(put);
    }

    @Override
    public void close() throws Exception {
        if (behaviorTable != null) behaviorTable.close();
        if (hbaseConn != null) hbaseConn.close();
    }
}

2.2 用户画像标签计算代码(Spark离线任务)

基于Spark SQL统计用户行为数据,生成“商品类目偏好”“返利敏感度”等标签,代码如下:

package cn.juwatech.tmallrebate.userportrait;

import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import static org.apache.spark.sql.functions.*;

/**
 * 用户画像标签离线计算(每日凌晨执行)
 */
public class UserPortraitCalculator {
    public static void main(String[] args) {
        // 1. 初始化SparkSession
        SparkSession spark = SparkSession.builder()
                .appName("User Portrait Calculator")
                .enableHiveSupport()
                .getOrCreate();

        // 2. 读取Hive中的用户行为数据(近30天)
        long thirtyDaysAgo = System.currentTimeMillis() - 30L * 24 * 60 * 60 * 1000;
        Dataset<Row> behaviorData = spark.sql(String.format(
                "SELECT user_id, product_id, behavior_type, rebate_rate, behavior_time " +
                "FROM tmall_behavior_hive " +
                "WHERE behavior_time >= %d", thirtyDaysAgo
        ));

        // 3. 计算“商品类目偏好”标签(统计用户浏览/购买最多的类目)
        Dataset<Row> categoryPreference = behaviorData
                // 关联商品表获取类目信息
                .join(spark.table("tmall_product_hive"), "product_id")
                // 过滤有效行为(浏览、购买)
                .filter(col("behavior_type").isin("VIEW", "PURCHASE"))
                // 按用户、类目分组计数
                .groupBy("user_id", "product_category")
                .agg(count("*").alias("behavior_count"))
                // 按用户分组,取计数最大的类目作为偏好标签
                .withColumn("rn", row_number().over(
                        org.apache.spark.sql.expressions.Window.partitionBy("user_id")
                                .orderBy(col("behavior_count").desc())
                ))
                .filter(col("rn").equalTo(1))
                .select(
                        col("user_id"),
                        col("product_category").alias("category_preference")
                );

        // 4. 计算“返利敏感度”标签(高/中/低)
        Dataset<Row> rebateSensitivity = behaviorData
                .filter(col("behavior_type").equalTo("PURCHASE"))
                .groupBy("user_id")
                // 计算平均返利比例,按阈值划分敏感度
                .agg(avg("rebate_rate").alias("avg_rebate_rate"))
                .withColumn("rebate_sensitivity", 
                        when(col("avg_rebate_rate").ge(0.15), "HIGH")
                        .when(col("avg_rebate_rate").ge(0.08), "MIDDLE")
                        .otherwise("LOW")
                )
                .select("user_id", "rebate_sensitivity");

        // 5. 合并标签写入HBase用户画像表
        categoryPreference.join(rebateSensitivity, "user_id")
                .write()
                .format("org.apache.hadoop.hbase.spark")
                .option("hbase.table", "tmall_user_portrait")
                .option("hbase.zookeeper.quorum", "zk-node1,zk-node2")
                .option("hbase.row.key", "user_id")
                .save();

        spark.stop();
    }
}

三、个性化推荐系统实现

3.1 推荐算法代码(协同过滤+标签匹配)

结合用户画像标签与协同过滤算法,生成个性化商品列表,代码如下:

package cn.juwatech.tmallrebate.recommend;

import cn.juwatech.tmallrebate.dto.UserPortraitDTO;
import cn.juwatech.tmallrebate.dto.ProductDTO;
import cn.juwatech.tmallrebate.mapper.ProductMapper;
import cn.juwatech.tmallrebate.service.UserPortraitService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.stream.Collectors;

/**
 * 个性化推荐服务
 */
@Service
public class PersonalRecommendService {

    @Autowired
    private UserPortraitService portraitService;
    @Autowired
    private ProductMapper productMapper;
    @Autowired
    private CollaborativeFilteringEngine cfEngine;

    /**
     * 生成用户个性化推荐列表(Top20)
     */
    public List<ProductDTO> getRecommendProducts(String userId) {
        // 1. 获取用户画像
        UserPortraitDTO portrait = portraitService.getPortraitByUserId(userId);
        if (portrait == null) {
            // 无画像数据时,返回热门商品
            return productMapper.getHotProducts(20);
        }

        // 2. 标签匹配:筛选用户偏好类目+高返利商品
        List<ProductDTO> tagMatchedProducts = productMapper.selectByConditions(
                portrait.getCategoryPreference(),
                portrait.getRebateSensitivity().equals("HIGH") ? 0.12 : 0.05,
                50  // 取前50个标签匹配商品
        );

        // 3. 协同过滤:基于相似用户行为推荐
        List<String> similarUserIds = cfEngine.getSimilarUsers(userId, 10);  // 获取10个相似用户
        List<ProductDTO> cfRecommendedProducts = productMapper.selectByUserIds(
                similarUserIds,
                50  // 取前50个协同过滤商品
        );

        // 4. 合并去重,按推荐分数排序(标签匹配权重0.6,协同过滤权重0.4)
        return tagMatchedProducts.stream()
                .map(product -> {
                    // 计算推荐分数
                    double tagScore = product.getRebateRate() >= 0.1 ? 0.8 : 0.6;
                    double cfScore = cfRecommendedProducts.contains(product) ? 0.4 : 0;
                    product.setRecommendScore(tagScore + cfScore);
                    return product;
                })
                .distinct()
                .sorted((p1, p2) -> Double.compare(p2.getRecommendScore(), p1.getRecommendScore()))
                .limit(20)  // 返回Top20
                .collect(Collectors.toList());
    }
}

// 协同过滤引擎(简化版)
class CollaborativeFilteringEngine {
    @Autowired
    private UserSimilarityMapper similarityMapper;

    /**
     * 获取相似用户ID列表
     */
    public List<String> getSimilarUsers(String userId, int limit) {
        // 查询用户相似度表,返回TopN相似用户
        return similarityMapper.selectSimilarUsers(userId, limit);
    }
}

3.2 推荐结果缓存与接口代码

将推荐结果缓存至Redis,减少计算开销,API接口代码如下:

package cn.juwatech.tmallrebate.controller;

import cn.juwatech.tmallrebate.dto.ProductDTO;
import cn.juwatech.tmallrebate.recommend.PersonalRecommendService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
import java.util.concurrent.TimeUnit;

@RestController
public class RecommendController {

    private static final String RECOMMEND_CACHE_KEY = "tmall:recommend:%s";
    private static final long CACHE_EXPIRE_HOURS = 2;

    @Autowired
    private PersonalRecommendService recommendService;
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @GetMapping("/api/recommend/personal")
    public List<ProductDTO> getPersonalRecommend(@RequestParam String userId) {
        // 1. 先查Redis缓存
        String cacheKey = String.format(RECOMMEND_CACHE_KEY, userId);
        List<ProductDTO> cachedRecommend = (List<ProductDTO>) redisTemplate.opsForValue().get(cacheKey);
        if (cachedRecommend != null && !cachedRecommend.isEmpty()) {
            return cachedRecommend;
        }

        // 2. 缓存未命中,调用推荐服务
        List<ProductDTO> recommendProducts = recommendService.getRecommendProducts(userId);

        // 3. 写入Redis缓存(设置2小时过期)
        redisTemplate.opsForValue().set(cacheKey, recommendProducts, CACHE_EXPIRE_HOURS, TimeUnit.HOURS);

        return recommendProducts;
    }
}

四、架构优化与性能调优

  1. 计算性能优化:将Spark离线任务拆分为增量计算(每日新增数据)与全量计算(每周一次),减少计算资源占用;
  2. 推荐延迟优化:通过Redis集群缓存推荐结果,热点用户推荐响应时间从500ms降至50ms以内;
  3. 数据倾斜处理:在Flink实时计算中,对用户ID进行哈希预分区,解决大用户行为数据倾斜问题。

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

Logo

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

更多推荐