我们日常生活中,经常会去网购,那么当你打开淘宝,拼多多等,从搜索中输入你想要的商品,开始浏览,加入购物车,下单结算等这一系列动作都可以作为用户行为的跟踪,可以给运营人员实时的分析当前用户的购买情况,也就可以相应的实时推荐一些商品.用户的购买行为可以有很多指标来衡量,比如:购买路径长度,购买次数,购买频率,感兴趣商品等.我们可以来看看购买路径长度这个指标,也就是从用户登录进App从开始到最终购买一系列动作的次数.比如,如果一个用户上来直接浏览了一次商品就添加进购物车,那么从这个购买路径长度说明该用户目的性很明确,那么推荐其他商品的意义可能不大.因此,我们可以定义一个规则,比方说购买路径长度大于某个动态配置的 阈值就触发推荐某些商品的操作,达到刺激消费的商业效果.
接下来,我们看看怎么用flink实现这个用户行为的跟踪.先简单画下处理的逻辑实现:
首先定义数据日志都会从Topic中生产出来,分为真实的用户行为事件源和可配置的数据源来进行消费;
在flink里面EventStream根据kafka特性有可能是多个分区的,所以需要将每个配置数据源都能与其每个分区数据做操作,所以可以通过把配置数据源做广播处理.
两个流数据再做connect操作,把结果sink下沉到想到的地方.具体使用union还是connect操作,需要看双流的类型,connect适合两条不同类型的流做操作.

在工程实现上,我们可以先构造两个流,用户事件流和配置流.
用户事件流格式为:
{"userId":"d8f3368aba5df27a39cbcfd36ce8084f","channel":"APP","eventType":"VIEW_PRODUCT","eventTime":"2019-10-24 09:27:11","data":{"productId":196}}配置流格式为:
{"channel":"APP","registerDate":"2019-01-01","historyPurchaseTimes":0,"maxPurchasePathLength":3}通过Kafka随机大量生产用户事件的流数据放在一个topic,配置流放在另外一个topic.
这里创建的实体模型包含,配置流的基本属性, 产品属性,用户事件属性,最终结果的属性.
下次再来实现把产生的流数据经过flink进行消费处理

文章转载自Nathan的笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




