导购返利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 研发团队转载请注明出处

相关新闻

【正则生成黄金标准】:IEEE最新白皮书认证的5大评估维度,92%开发者从未用过!

【正则生成黄金标准】:IEEE最新白皮书认证的5大评估维度,92%开发者从未用过!

更多请点击: https://kaifayun.com 第一章:AI 生成正则表达式的范式革命与黄金标准定义 传统正则表达式编写长期依赖人工经验与反复调试,而AI驱动的正则生成正从根本上重构这一范式:从“人写规则”转向“人描述意图,A…

2026/7/24 19:12:22阅读更多 →
ERC-725 与 ERC-735 去中心化身份实现:声明发布、验证请求与链上凭证管理

ERC-725 与 ERC-735 去中心化身份实现:声明发布、验证请求与链上凭证管理

ERC-725 与 ERC-735 去中心化身份实现:声明发布、验证请求与链上凭证管理 一、DID 的标准不止一种,但 725735 的组合最完整 ERC-725 和 ERC-735 是 Ethereum 上实现去中心化身份的两项核心标准。ERC-725 定义链上身份的存储结构和权限管理:…

2026/7/24 19:12:22阅读更多 →
AI生成社交媒体封面:凌晨2点还在改稿?用这5个自动化工作流,日均产出47张合规封面

AI生成社交媒体封面:凌晨2点还在改稿?用这5个自动化工作流,日均产出47张合规封面

更多请点击: https://kaifayun.com 第一章:AI生成社交媒体封面 AI生成社交媒体封面正迅速成为数字内容创作者的核心工作流之一。借助多模态大模型与扩散模型技术,用户仅需输入简洁的文本提示(prompt),即可…

2026/7/24 19:12:22阅读更多 →
解决广色域显示器过饱和问题:novideo_srgb色彩校准终极指南 [特殊字符]

解决广色域显示器过饱和问题:novideo_srgb色彩校准终极指南 [特殊字符]

解决广色域显示器过饱和问题:novideo_srgb色彩校准终极指南 🎨 【免费下载链接】novideo_srgb Calibrate monitors to sRGB or other color spaces on NVIDIA GPUs, based on EDID data or ICC profiles 项目地址: https://gitcode.com/gh_mirrors/no/…

2026/7/24 20:48:40阅读更多 →
机器视觉8 —— CogPatInspectTool 缺陷对比工具(加画轮廓)

机器视觉8 —— CogPatInspectTool 缺陷对比工具(加画轮廓)

CogPatInspectTool主要用于缺陷检测和复杂模式分析功能特点模式检测:能够根据图像中的特征和模式来检测目标物体,即使在复杂背景下也能准确识位置和角度检测:可以检测并识别目标物体的位置和角度信息,为后续的分析和处理提供基多目…

2026/7/24 20:48:40阅读更多 →
3分钟极速指南:Deepin Boot Maker启动盘制作终极教程

3分钟极速指南:Deepin Boot Maker启动盘制作终极教程

3分钟极速指南:Deepin Boot Maker启动盘制作终极教程 【免费下载链接】deepin-boot-maker 项目地址: https://gitcode.com/gh_mirrors/de/deepin-boot-maker 你是否曾经为制作系统启动盘而烦恼?Deepin Boot Maker作为一款轻量级启动盘制作工具&a…

2026/7/24 20:48:39阅读更多 →
springboot校园在线拍卖系统13216--计算机设计/毕业设计

springboot校园在线拍卖系统13216--计算机设计/毕业设计

前言 📌博主介绍:一线全栈工程师,毕设实战引路人。技术栈覆盖Java、Python、C#、PHP、Node.js及UniApp跨端开发,擅长多语言项目落地与架构设计。持续分享毕设源码、开题报告、技术选型心得与职场踩坑经验。用工程化思维写代码&am…

2026/7/24 20:48:39阅读更多 →
DaVinci CFG的BSWM配置

DaVinci CFG的BSWM配置

General Settings ErrorDetect错误检测 MainFunctionPeriod主函数调度周期 ModeCheck将检查传递给模式请求API的模式值 UserConfigurationFile头文件 VersionInfoApi使能版本接口 DataTypeMappingSetRef引用DEV里面的Type Mapping Set,这张表是自

2026/7/24 20:48:39阅读更多 →
Locale Emulator终极指南:快速解决多语言软件乱码问题

Locale Emulator终极指南:快速解决多语言软件乱码问题

Locale Emulator终极指南:快速解决多语言软件乱码问题 【免费下载链接】Locale-Emulator Yet Another System Region and Language Simulator 项目地址: https://gitcode.com/gh_mirrors/lo/Locale-Emulator 你是否遇到过下载日文游戏或软件时,打…

2026/7/24 20:46:39阅读更多 →
Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/24 0:58:53阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/24 0:58:53阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/24 0:58:53阅读更多 →
我的编程之路:第一篇博客

我的编程之路:第一篇博客

大家好,我是一名编程初学者,同时这也是我编程学习之路上的第一篇博客。在这里,我想要向大家介绍我的一些想法和规划。a.自我介绍我是一个刚刚接触编程的新手,目前在学习c语言,我对编程世界充满了强烈的好奇。当然&…

2026/7/24 0:00:06阅读更多 →
【LeetCode 54】螺旋矩阵

【LeetCode 54】螺旋矩阵

问题描述: 解法: 1、模拟(参考自【LeetCode 54】螺旋矩阵-CSDN博客) int *spiralOrder(int **matrix, int matrixSize, int *matrixColSize, int *returnSize) {static const int dirs[4][2] {{0, 1}, {1, 0}, {0, -1}, {-1, …

2026/7/24 0:00:06阅读更多 →
2026 WAIC:模型隐身、智能体疯野,厂商竞赛聚焦办公场景与商业闭环

2026 WAIC:模型隐身、智能体疯野,厂商竞赛聚焦办公场景与商业闭环

知春路不相信模型领先今年WAIC大会,昔日AI六小龙来了五家,分别是Kimi、阶跃星辰、Minimax、百川智能、零一万物。连放弃基模的百川和零一万物都来了,唯一缺席的竟是近几个月来风光无限的智谱。(DeepSeek一直不参加)WAI…

2026/7/24 0:00:06阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/23 22:58:43阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/24 19:00:40阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/24 19:00:40阅读更多 →