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