导购返利APP用户行为日志采集与实时返利计算的流式处理架构
大家好,我是省赚客APP研发者微赚淘客!
在导购返利业务中,订单追踪的实时性与准确性是核心竞争力。传统的T+1离线批处理模式已无法满足用户对“下单即见返利”的体验期待。为此,我们构建了基于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{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 研发团队,转载请注明出处!