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

发布时间:2026/9/12 8:19:56

导购返利APP用户行为日志采集与实时返利计算的流式处理架构 导购返利APP用户行为日志采集与实时返利计算的流式处理架构大家好我是省赚客APP研发者微赚淘客在导购返利业务中订单追踪的实时性与准确性是核心竞争力。传统的T1离线批处理模式已无法满足用户对“下单即见返利”的体验期待。为此我们构建了基于Apache Flink的实时流式处理架构实现了从用户行为采集到返利金额计算的毫秒级响应。一、 整体架构设计我们的实时返利计算系统遵循经典的Lambda架构思想但侧重于速度层Speed Layer的实时处理能力。整体数据流如下数据采集层APP端用户行为点击、下单通过SDK上报至Nginx再由Filebeat采集写入Kafka。消息队列层Kafka作为高吞吐的日志缓冲解耦数据生产与消费。流式计算层Flink消费Kafka数据进行ETL、订单匹配、返利计算。结果存储层计算结果写入Redis供APP实时查询和MySQL持久化。二、 用户行为日志采集首先我们需要定义统一的用户行为日志格式以便下游系统解析。1. 日志数据模型 (Java POJO)packagejuwatech.cn.tracker.model;importjava.io.Serializable;/** * 用户行为日志实体 * author juwatech.cn */publicclassUserActionLogimplementsSerializable{privatestaticfinallongserialVersionUID1L;// 用户IDprivateStringuserId;// 行为类型: CLICK, ORDER, PAYprivateStringactionType;// 商品IDprivateStringitemId;// 订单ID (下单行为时有值)privateStringorderId;// 订单金额privateDoubleorderAmount;// 时间戳privateLongtimestamp;// 渠道来源 (淘宝/京东/拼多多)privateStringchannel;// Getters and SetterspublicStringgetUserId(){returnuserId;}publicvoidsetUserId(StringuserId){this.userIduserId;}publicStringgetActionType(){returnactionType;}publicvoidsetActionType(StringactionType){this.actionTypeactionType;}publicStringgetItemId(){returnitemId;}publicvoidsetItemId(StringitemId){this.itemIditemId;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderIdorderId;}publicDoublegetOrderAmount(){returnorderAmount;}publicvoidsetOrderAmount(DoubleorderAmount){this.orderAmountorderAmount;}publicLonggetTimestamp(){returntimestamp;}publicvoidsetTimestamp(Longtimestamp){this.timestamptimestamp;}publicStringgetChannel(){returnchannel;}publicvoidsetChannel(Stringchannel){this.channelchannel;}}2. 日志采集SDK (Android端伪代码)packagejuwatech.cn.tracker.sdk;importandroid.content.Context;importandroid.os.AsyncTask;importorg.json.JSONObject;/** * 埋点SDK核心类 * author juwatech.cn */publicclassTrackerSDK{privatestaticfinalStringSERVER_URLhttps://log.juwatech.cn/collect;privateContextcontext;publicTrackerSDK(Contextcontext){this.contextcontext;}/** * 上报用户行为 */publicvoidtrack(StringactionType,StringitemId,StringorderId,doubleamount){newUploadTask().execute(actionType,itemId,orderId,String.valueOf(amount));}privateclassUploadTaskextendsAsyncTaskString,Void,Void{OverrideprotectedVoiddoInBackground(String...params){try{JSONObjectjsonnewJSONObject();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(){// 获取设备唯一标识returndevice_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. 获取执行环境finalStreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(4);// 2. 配置Kafka消费者PropertiespropertiesnewProperties();properties.setProperty(bootstrap.servers,localhost:9092);properties.setProperty(group.id,rebate-consumer-group);FlinkKafkaConsumerStringkafkaSourcenewFlinkKafkaConsumer(user-action-topic,newSimpleStringSchema(),properties);// 3. 添加数据源DataStreamStringrawStreamenv.addSource(kafkaSource);// 4. 数据转换JSON字符串 - UserActionLog对象DataStreamUserActionLoglogStreamrawStream.map(json-JSON.parseObject(json,UserActionLog.class));// 5. 过滤出下单行为DataStreamUserActionLogorderStreamlogStream.filter(log-ORDER.equals(log.getActionType()));// 6. 核心计算计算返利金额DataStreamRebateResultresultStreamorderStream.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 */publicclassRebateCalculationFunctionimplementsMapFunctionUserActionLog,RebateResult{OverridepublicRebateResultmap(UserActionLoglog)throwsException{RebateResultresultnewRebateResult();result.setUserId(log.getUserId());result.setOrderId(log.getOrderId());result.setItemId(log.getItemId());// 模拟返利比例查询 (实际应查询维表或缓存)doublerebateRategetRebateRate(log.getChannel(),log.getItemId());// 计算返利金额doublerebateAmountlog.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.userIduserId;}publicStringgetOrderId(){returnorderId;}publicvoidsetOrderId(StringorderId){this.orderIdorderId;}publicStringgetItemId(){returnitemId;}publicvoidsetItemId(StringitemId){this.itemIditemId;}publicDoublegetRebateAmount(){returnrebateAmount;}publicvoidsetRebateAmount(DoublerebateAmount){this.rebateAmountrebateAmount;}publicLonggetCalcTime(){returncalcTime;}publicvoidsetCalcTime(LongcalcTime){this.calcTimecalcTime;}}4. 自定义Sink写入Redispackagejuwatech.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 */publicclassCustomRedisSinkextendsRedisSinkRebateResult{publicCustomRedisSink(){super(newRedisConnectionConfig(localhost,6379),newRebateRedisMapper());}privatestaticclassRebateRedisMapperimplementsRedisMapperRebateResult{OverridepublicRedisCommandDescriptiongetCommandDescription(){// 使用HASH结构存储: keyrebate:userId, fieldorderId, valueamountreturnnewRedisCommandDescription(RedisCommand.HSET,rebate:);}OverridepublicStringgetKeyFromData(RebateResultdata){returndata.getUserId();}OverridepublicStringgetValueFromData(RebateResultdata){returndata.getOrderId():data.getRebateAmount();}}}通过这套流式处理架构我们将返利到账时间从小时级缩短到了秒级。当用户在省赚客APP下单后Flink任务几乎实时捕获订单日志完成返利计算并更新Redis用户刷新页面即可看到预计返利金额极大地提升了用户粘性与信任度。本文著作权归 省赚客app 研发团队转载请注明出处
延伸阅读

更多相关文章

2026/9/12 8:15:09

CPU、内存、硬盘:电脑硬件三件套的分工与协作原理

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 8:15:09

SNMP协议栈选型:Net-SNMP与国产自研在信创迁移中的实战对比

SNMP这东西,平时不显山不露水,一旦碰上设备纳管、网络监控、平台对接,它就是绕不过去的那道坎。这个标题里几个关键词——SNMP协议栈、SNMP SDK、Net-SNMP、国产自研、信创,我基本都折腾过,尤其最近两年做信创环境下的…

2026/9/12 8:15:09

基于Matlab的手指手掌静脉识别实现与算法详解

简介:面向机器视觉课程创新实践,这份手指手掌静脉识别Matlab工程聚焦手部静脉图像预处理算法实验研究。项目完整覆盖静脉识别链路,对手指和手掌分别进行轮廓分割、感兴趣区域(ROI)截取、静脉纹理增强与分割&#xff0c…

2026/9/12 8:15:09

制定项目章程全解析:启动过程组考点与10道真题拆解

很多备考信息系统项目管理师的同学,看到“启动过程组”就开始头疼:不就是一个立项流程吗,怎么还整出输入、工具与输出了?结果做题时各种翻车——商业文件被当成项目文件、项目章程以为是项目经理发布的、启动大会当成正式过程………

2026/9/12 8:15:08

分数阶模型与FOTF工具箱:从传递函数到分数阶PID整定实战

简介:专门处理分数阶系统建模与仿真的MATLAB微积分工具箱,面向科研人员、工程师及高校学生,解决分数阶微分方程求解、传递函数构建与数值积分等核心问题。压缩包共289个文件,总大小约2.32MB,含164个m脚本、53个mdl与37…

2026/9/12 8:10:08

GEO产业图谱解析与企业竞争力分析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/12 2:05:33

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/12 3:55:12

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/12 0:04:17

MATLAB仿生优化框架:长鼻浣熊算法多策略融合实现

简介:本资源是一份面向智能优化算法研究者与MATLAB初学者的仿生智能算法实践代码包,聚焦于长鼻浣熊优化算法(COA)的多策略改进与性能验证。针对传统COA易陷局部最优、收敛精度不足等问题,作者融合Circle映射初始化提升…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 JavaWeb 的校园一卡通管理系统的设计与实现 基于 JavaWeb 的校园卡业务管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 Java 的图书馆借阅管理平台的搭建与实现 基于 Java 的图书馆综合管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 6:29:36

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/10 15:19:50

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/12 6:37:43

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码