天猫返利app的大数据处理架构:用户画像与个性化推荐系统设计
·
天猫返利app的大数据处理架构:用户画像与个性化推荐系统设计
大家好,我是阿可,微赚淘客系统及省赚客APP创始人,是个冬天不穿秋裤,天冷也要风度的程序猿!
在天猫返利app业务中,用户画像与个性化推荐是核心竞争力——通过分析用户浏览、下单、返利领取等行为数据,构建精准用户画像,再基于画像推送高匹配度的天猫商品及返利活动,可将转化率提升30%以上。为此,我们设计了“实时+离线”双轨大数据处理架构,基于Hadoop生态与实时计算引擎,实现从数据采集到推荐落地的全链路闭环。以下从架构分层、用户画像构建、个性化推荐实现三方面展开,附关键代码示例。

一、大数据处理架构整体设计
架构采用五层分布式设计,各层职责与技术选型如下:
- 数据采集层:通过Flume采集用户行为日志(如商品点击、加入购物车),使用Canal同步天猫开放平台API数据(商品价格、返利比例),数据统一写入Kafka消息队列;
- 数据存储层:HDFS存储离线原始数据,HBase存储用户画像标签(支持高并发读写),Redis缓存实时推荐结果;
- 计算层:离线计算用Spark SQL处理用户行为数据生成画像标签,实时计算用Flink处理用户实时行为(如当前浏览商品)触发动态推荐;
- 画像层:构建“用户基础属性+行为偏好+返利敏感度”三维标签体系,通过标签权重算法生成用户画像;
- 推荐层:基于协同过滤与标签匹配算法,生成个性化商品推荐列表,通过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;
}
}
四、架构优化与性能调优
- 计算性能优化:将Spark离线任务拆分为增量计算(每日新增数据)与全量计算(每周一次),减少计算资源占用;
- 推荐延迟优化:通过Redis集群缓存推荐结果,热点用户推荐响应时间从500ms降至50ms以内;
- 数据倾斜处理:在Flink实时计算中,对用户ID进行哈希预分区,解决大用户行为数据倾斜问题。
本文著作权归聚娃科技省赚客app开发者团队,转载请注明出处!
更多推荐



所有评论(0)