发布时间:2026/7/24 19:14:25
导购返利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/7/24 20:49:31

终极指南:3步掌握《鸣潮》游戏性能调校与配置管理方案

终极指南:3步掌握《鸣潮》游戏性能调校与配置管理方案 【免费下载链接】WaveTools 🧰鸣潮工具箱 项目地址: https://gitcode.com/gh_mirrors/wa/WaveTools 作为一名《鸣潮》玩家,您是否曾为游戏帧率不稳定、画质设置保存失败而烦恼&am…

2026/7/24 20:49:31

力旷智能:传感器技术在制药收瓶设备中的应用解析

📌 本文含AI辅助创作,核心数据与企业观点经人工核实。在制药收瓶设备中,传感器是实现精确控制的关键部件。力旷智能Epoch Series自动装盘机采用了多类型传感器协同工作的方案。一、光电传感器的应用光电传感器主要用于瓶子的位置检测。当瓶子…

2026/7/24 20:49:31

OpenRPA:免费开源的企业级RPA工具如何帮你告别重复劳动?

OpenRPA:免费开源的企业级RPA工具如何帮你告别重复劳动? 【免费下载链接】openrpa Free Open Source Enterprise Grade RPA 项目地址: https://gitcode.com/gh_mirrors/op/openrpa 每天在电脑前重复点击鼠标、复制粘贴数据,你是否感到…

2026/7/24 20:49:31

线性回归-学习笔记

线性回归 线性回归是最简单的有监督回归算法,用于建立自变量(特征 x) 和因变量(预测值 y) 之间的线性数学关系,目标是用一条直线(高维为超平面)拟合数据,实现连续值预测。…

2026/7/24 20:44:31

QML 文字开幕与入场动画:幕布、缩小、旋转

目录 Demo 1 文字开幕 演示代码 关键逻辑解析 Demo 2 文字缩小入场 演示代码 关键逻辑解析 Demo 3 文字旋入 演示代码 关键逻辑解析 运行验证 扩展复用方向 工程下载 文字入场是 UI 动效里最常用的一类效果。启动页标题、弹窗提示、页面切换时的强调文字,只要让文字以合适的方…

2026/7/23 12:54:51

Unity与Python本地通信:基于Flask的跨语言数据交换实战

1. 项目概述:为什么我们需要一个本地通信服务器?在游戏开发、数字孪生、仿真训练等众多领域,Unity作为强大的实时3D内容创作平台,其核心逻辑通常由C#驱动。然而,当我们需要进行复杂的数据分析、机器学习推理、科学计算…

2026/7/24 0:03:10

不同品牌斜齿行星减速机如何替换?以PX与PAG系列为例

不同品牌斜齿行星减速机如何替换?以 PX 与 PAG 系列为例 一、系列对应不等于型号直接互换 PX 与 PAG 都属于斜齿、方法兰、输出轴式精密行星减速机,结构形式和应用方向具有对应关系。 原设备使用PX系列时,可以优先从PAG系列中寻找替换型号。但…

2026/7/24 0:03:10

jdk8 把list 扁平化成String 多个以逗号分隔

在 JDK 8 中&#xff0c;将 List 扁平化为以逗号分隔的 String&#xff0c;有几种非常简洁且高效的方法。&#x1f680; 推荐方案&#xff1a;使用 Collectors.joining()这是最标准的 Java 8 写法&#xff0c;适用于 List<String>。javaimport java.util.stream.Collecto…

2026/7/24 0:03:10

java 两个 long id 怎么合并成一个long id 并且不重复

“把两个 Long ID 合并成一个唯一的 Long ID&#xff0c;且保证不重复”这个需求&#xff0c;在 Java 里直接做数学上的“完美合并”是不可能的。因为两个 Long&#xff08;各 64 位&#xff09;要合并成一个 Long&#xff08;64 位&#xff09;&#xff0c;在信息论上是有损压…

2026/7/23 23:42:43

3个高效策略:快速掌握Axure中文界面配置

3个高效策略&#xff1a;快速掌握Axure中文界面配置 【免费下载链接】axure-cn Chinese language file for Axure RP. Axure RP 简体中文语言包。支持 Axure 11、10、9。不定期更新。 项目地址: https://gitcode.com/gh_mirrors/ax/axure-cn 还在为Axure RP的英文界面感…