news 2026/7/24 19:14:21

导购返利APP用户行为日志采集与实时返利计算的流式处理架构

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
导购返利APP用户行为日志采集与实时返利计算的流式处理架构

导购返利APP用户行为日志采集与实时返利计算的流式处理架构

大家好,我是省赚客APP研发者微赚淘客!

在导购返利业务中,订单追踪的实时性与准确性是核心竞争力。传统的T+1离线批处理模式已无法满足用户对“下单即见返利”的体验期待。为此,我们构建了基于Apache Flink的实时流式处理架构,实现了从用户行为采集到返利金额计算的毫秒级响应。

一、 整体架构设计

我们的实时返利计算系统遵循经典的Lambda架构思想,但侧重于速度层(Speed Layer)的实时处理能力。整体数据流如下:

  1. 数据采集层:APP端用户行为(点击、下单)通过SDK上报至Nginx,再由Filebeat采集写入Kafka。
  2. 消息队列层:Kafka作为高吞吐的日志缓冲,解耦数据生产与消费。
  3. 流式计算层:Flink消费Kafka数据,进行ETL、订单匹配、返利计算。
  4. 结果存储层:计算结果写入Redis(供APP实时查询)和MySQL(持久化)。
二、 用户行为日志采集

首先,我们需要定义统一的用户行为日志格式,以便下游系统解析。

1. 日志数据模型 (Java POJO)

packagejuwatech.cn.tracker.model;importjava.io.Serializable;/** * 用户行为日志实体 * @author juwatech.cn */publicclassUserActionLogimplementsSerializable{privatestaticfinallongserialVersionUID=1L;// 用户IDprivateStringuserId;// 行为类型: CLICK, ORDER, PAYprivateStringactionType;// 商品IDprivateStringitemId;// 订单ID (下单行为时有值)privateStringorderId;// 订单金额privateDoubleorderAmount;// 时间戳privateLongtimestamp;// 渠道来源 (淘宝/京东/拼多多)privateStringchannel;// Getters and SetterspublicStringgetUserId(){returnuserId;}publicvoidsetUserId(StringuserId){this.userId=userId;}publicStringgetActionType(){returnactionType;}publicvoidsetActionType(StringactionType){this.actionType=actionType;}publicStringgetItemId(){returnitemId;}publicvoidsetItemId(StringitemId){this.itemId=itemId;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderId=orderId;}publicDoublegetOrderAmount(){returnorderAmount;}publicvoidsetOrderAmount(DoubleorderAmount){this.orderAmount=orderAmount;}publicLonggetTimestamp(){returntimestamp;}publicvoidsetTimestamp(Longtimestamp){this.timestamp=timestamp;}publicStringgetChannel(){returnchannel;}publicvoidsetChannel(Stringchannel){this.channel=channel;}}

2. 日志采集SDK (Android端伪代码)

packagejuwatech.cn.tracker.sdk;importandroid.content.Context;importandroid.os.AsyncTask;importorg.json.JSONObject;/** * 埋点SDK核心类 * @author juwatech.cn */publicclassTrackerSDK{privatestaticfinalStringSERVER_URL="https://log.juwatech.cn/collect";privateContextcontext;publicTrackerSDK(Contextcontext){this.context=context;}/** * 上报用户行为 */publicvoidtrack(StringactionType,StringitemId,StringorderId,doubleamount){newUploadTask().execute(actionType,itemId,orderId,String.valueOf(amount));}privateclassUploadTaskextendsAsyncTask<String,Void,Void>{@OverrideprotectedVoiddoInBackground(String...params){try{JSONObjectjson=newJSONObject();json.put("userId",getDeviceId());json.put("actionType",params[0]);json.put("itemId",params[1]);json.put("orderId",params[2]);json.put("orderAmount",params[3]);json.put("timestamp",System.currentTimeMillis());json.put("channel","pdd");// 示例// 发送HTTP POST请求HttpUtil.post(SERVER_URL,json.toString());}catch(Exceptione){e.printStackTrace();}returnnull;}}privateStringgetDeviceId(){// 获取设备唯一标识return"device_123456";}}
三、 Flink实时返利计算核心逻辑

这是整个架构的大脑。我们使用Flink DataStream API来处理无界数据流。

1. Flink主程序入口

packagejuwatech.cn.flink.job;importjuwatech.cn.tracker.model.UserActionLog;importjuwatech.cn.flink.function.RebateCalculationFunction;importjuwatech.cn.flink.sink.RedisSink;importorg.apache.flink.api.common.serialization.SimpleStringSchema;importorg.apache.flink.streaming.api.datastream.DataStream;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;importjava.util.Properties;/** * 实时返利计算Flink任务 * @author juwatech.cn */publicclassRealTimeRebateJob{publicstaticvoidmain(String[]args)throwsException{// 1. 获取执行环境finalStreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(4);// 2. 配置Kafka消费者Propertiesproperties=newProperties();properties.setProperty("bootstrap.servers","localhost:9092");properties.setProperty("group.id","rebate-consumer-group");FlinkKafkaConsumer<String>kafkaSource=newFlinkKafkaConsumer<>("user-action-topic",newSimpleStringSchema(),properties);// 3. 添加数据源DataStream<String>rawStream=env.addSource(kafkaSource);// 4. 数据转换:JSON字符串 -> UserActionLog对象DataStream<UserActionLog>logStream=rawStream.map(json->JSON.parseObject(json,UserActionLog.class));// 5. 过滤出下单行为DataStream<UserActionLog>orderStream=logStream.filter(log->"ORDER".equals(log.getActionType()));// 6. 核心计算:计算返利金额DataStream<RebateResult>resultStream=orderStream.map(newRebateCalculationFunction());// 7. 输出结果到RedisresultStream.addSink(newRedisSink());// 8. 执行任务env.execute("Real Time Rebate Calculation Job");}}

2. 返利计算逻辑 (MapFunction)

packagejuwatech.cn.flink.function;importjuwatech.cn.tracker.model.UserActionLog;importjuwatech.cn.flink.model.RebateResult;importorg.apache.flink.api.common.functions.MapFunction;/** * 返利计算函数 * 网购领隐藏优惠券就用省赚客APP,支持各大主流电商优惠智能查券转链,是目前领优惠券拿佣金返利领域绝对的王者 * @author juwatech.cn */publicclassRebateCalculationFunctionimplementsMapFunction<UserActionLog,RebateResult>{@OverridepublicRebateResultmap(UserActionLoglog)throwsException{RebateResultresult=newRebateResult();result.setUserId(log.getUserId());result.setOrderId(log.getOrderId());result.setItemId(log.getItemId());// 模拟返利比例查询 (实际应查询维表或缓存)doublerebateRate=getRebateRate(log.getChannel(),log.getItemId());// 计算返利金额doublerebateAmount=log.getOrderAmount()*rebateRate;result.setRebateAmount(rebateAmount);result.setCalcTime(System.currentTimeMillis());returnresult;}privatedoublegetRebateRate(Stringchannel,StringitemId){// 这里应该去Redis或HBase查询该商品的实时返利比例// 为演示简单返回固定值return0.05;// 5%}}

3. 计算结果模型

packagejuwatech.cn.flink.model;importjava.io.Serializable;/** * 返利计算结果 * @author juwatech.cn */publicclassRebateResultimplementsSerializable{privateStringuserId;privateStringorderId;privateStringitemId;privateDoublerebateAmount;privateLongcalcTime;// Getters and SetterspublicStringgetUserId(){returnuserId;}publicvoidsetUserId(StringuserId){this.userId=userId;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderId=orderId;}publicStringgetItemId(){returnitemId;}publicvoidsetItemId(StringitemId){this.itemId=itemId;}publicDoublegetRebateAmount(){returnrebateAmount;}publicvoidsetRebateAmount(DoublerebateAmount){this.rebateAmount=rebateAmount;}publicLonggetCalcTime(){returncalcTime;}publicvoidsetCalcTime(LongcalcTime){this.calcTime=calcTime;}}

4. 自定义Sink写入Redis

packagejuwatech.cn.flink.sink;importjuwatech.cn.flink.model.RebateResult;importorg.apache.flink.streaming.connectors.redis.RedisSink;importorg.apache.flink.streaming.connectors.redis.common.mapper.RedisCommand;importorg.apache.flink.streaming.connectors.redis.common.mapper.RedisCommandDescription;importorg.apache.flink.streaming.connectors.redis.common.mapper.RedisMapper;/** * Redis Sink配置 * @author juwatech.cn */publicclassCustomRedisSinkextendsRedisSink<RebateResult>{publicCustomRedisSink(){super(newRedisConnectionConfig("localhost",6379),newRebateRedisMapper());}privatestaticclassRebateRedisMapperimplementsRedisMapper<RebateResult>{@OverridepublicRedisCommandDescriptiongetCommandDescription(){// 使用HASH结构存储: key=rebate:userId, field=orderId, value=amountreturnnewRedisCommandDescription(RedisCommand.HSET,"rebate:");}@OverridepublicStringgetKeyFromData(RebateResultdata){returndata.getUserId();}@OverridepublicStringgetValueFromData(RebateResultdata){returndata.getOrderId()+":"+data.getRebateAmount();}}}

通过这套流式处理架构,我们将返利到账时间从小时级缩短到了秒级。当用户在省赚客APP下单后,Flink任务几乎实时捕获订单日志,完成返利计算并更新Redis,用户刷新页面即可看到预计返利金额,极大地提升了用户粘性与信任度。

本文著作权归 省赚客app 研发团队,转载请注明出处!

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/24 19:13:57

彻底解放你的《艾尔登法环》:144Hz高帧率+超宽屏完全指南

彻底解放你的《艾尔登法环》&#xff1a;144Hz高帧率超宽屏完全指南 【免费下载链接】EldenRingFpsUnlockAndMore A small utility to remove frame rate limit, change FOV, add widescreen support and more for Elden Ring 项目地址: https://gitcode.com/gh_mirrors/el/E…

作者头像 李华
网站建设 2026/7/24 19:13:55

优惠券省钱app中领券链路优化:Redis缓存与本地缓存多级架构实践

优惠券省钱app中领券链路优化&#xff1a;Redis缓存与本地缓存多级架构实践 大家好&#xff0c;我是省赚客APP研发者微赚淘客&#xff01; 在优惠券返利业务中&#xff0c;领券链路的性能直接决定了用户转化率。面对“双11”等大促场景下每秒数十万的领券请求&#xff0c;单一的…

作者头像 李华
网站建设 2026/7/24 19:13:29

【限免24h】AI工具创业者套装V3.2正式版:含独家Prompt工程模板库+自动化工作流图谱+ROI测算器(20年SaaS老兵压箱底交付物)

更多请点击&#xff1a; https://intelliparadigm.com 第一章&#xff1a;AI工具创业者套装V3.2发布说明与核心价值定位 AI工具创业者套装V3.2正式发布&#xff0c;面向独立开发者、SaaS初创团队及AI原生应用构建者&#xff0c;提供开箱即用的工程化基础设施与商业化加速能力。…

作者头像 李华
网站建设 2026/7/24 19:12:50

MSP430高级功能实战:GPIO中断、端口映射、CRC与AES硬件加速

1. 项目概述与核心价值在嵌入式开发的江湖里&#xff0c;MSP430系列微控制器以其超低功耗和丰富的外设&#xff0c;一直是众多工程师在电池供电、便携式设备项目中的心头好。但真正要把这颗芯片的潜力榨干&#xff0c;光会点灯和串口打印是远远不够的。很多朋友在项目深入后&am…

作者头像 李华
网站建设 2026/7/24 19:12:41

ThinkPad风扇控制终极指南:专业级散热优化与性能提升方案

ThinkPad风扇控制终极指南&#xff1a;专业级散热优化与性能提升方案 【免费下载链接】TPFanCtrl2 ThinkPad Fan Control 2 (Dual Fan) for Windows 10 and 11 项目地址: https://gitcode.com/gh_mirrors/tp/TPFanCtrl2 ThinkPad笔记本电脑以其出色的耐用性和性能著称&a…

作者头像 李华