一、全体的なアプローチ

二、コードの解析
1)LRデータの準備:
1、データを結合する。ユーザーが過去に見た商品は、行動に基づいて好む0または好まない1で区別し、ユーザーが見たことがない商品は2でマークします。
// ユーザーが商品に好むかどうかを判断します。ユーザーが注文したり、カートに入れるなら好むとします。
val isLove: UserDefinedFunction = udf{
(act:String)=>{
if(act.equalsIgnoreCase("BROWSE")
||act.equalsIgnoreCase("COLLECT")){
0
}else{
1
}
}
}
import spark.implicits._
// グローバルホットセールデータの取得
// (cust_id,good_id,rank)
val hot = HDFSConnection.readDataToHDFS(spark,"/myshops/dwd_hotsell")
.select($"cust_id",$"good_id")
// グループ再現データの取得
val group = HDFSConnection.readDataToHDFS(spark,"/myshops/dwd_kMeans")
.select($"cust_id",$"good_id")
// ALS再現データの取得
val als = HDFSConnection.readDataToHDFS(spark,"/myshops/dwd_ALS_Iter20")
.select($"cust_id",$"good_id")
// 注文データの取得、ユーザーが注文したりカートに入れるなら好む、それ以外は好まない
val order = spark.sparkContext
.textFile("file:///D:/logs/virtualLogs/*.log")
.map(line=>{
val arr = line.split(" ")
(arr(0),arr(2),arr(3))
})
.toDF("act","cust_id","good_id")
.withColumn("flag",isLove($"act"))
.drop("act")
.distinct()
.cache()
// 三路再現結果の結合(新規ユーザーには2を追加)
// 新規ユーザーが全く見たことのない商品には2を追加
val all = hot.union(group).union(als)
.join(order,Seq("cust_id","good_id"),"left")
.na.fill(2)
2、LRモデルが必要とするデータを準備します:ラベル:好むか否か、特徴:ユーザーと商品の属性、正規化
// 簡単なデータの正規化
val priceNormalize: UserDefinedFunction =udf{
(price:String)=>{
// maxscale & minscale
val p:Double = price.toDouble
p/(10000+p)
}
}
def goodNumberFormat(spark: SparkSession): DataFrame ={
val good_infos = MYSQLConnection.readMySql(spark,"goods")
.filter("is_sale=1")
.drop("spu_pro_name","tags","content","good_name","created_at","update_at","good_img_pos","sku_good_code")
// ブランドの数値化処理
val brand_index = new StringIndexer().setInputCol("brand_name").setOutputCol("brand")
val bi = brand_index.fit(good_infos).transform(good_infos)
// 商品カテゴリの数値化
val type_index = new StringIndexer().setInputCol("cate_name").setOutputCol("cate")
val ct = type_index.fit(bi).transform(bi)
// 原価と現在価格の正規化
import spark.implicits._
val pc = ct.withColumn("nprice",priceNormalize($"price"))
.withColumn("noriginal",priceNormalize($"original"))
.withColumn("nsku_num",priceNormalize($"sku_num"))
.drop("price","original","sku_num")
// 特徴値を数値化
val feat_index = new StringIndexer().setInputCol("spu_pro_value").setOutputCol("pro_value")
feat_index.fit(pc).transform(pc).drop("spu_pro_value")
}
// 各列にLR回帰アルゴリズムが必要とするユーザーの自然属性、ユーザーの行動属性、商品の自然属性を追加
val user_info_df = KMeansHandler.user_act_info(spark)
// データベースから商品に影響を与える自然属性を取得
val good_infos = goodNumberFormat(spark)
// 3路再現結果とユーザー情報および商品情報を関連付ける
val ddf = all.join(user_info_df,Seq("cust_id"),"inner")
.join(good_infos,Seq("good_id"),"inner")
// 全体のデータをDouble型に変換
val columns = ddf.columns.map(f => col(f).cast(DoubleType))
val num_fmt = ddf.select(columns:_*)
// 特徴列を集約して密度ベクトルを作成
val va = new VectorAssembler().setInputCols(
Array("province_id","city_id","district_id","sex","marital_status","education_id","vocation","post","compId","mslevel","reg_date","lasttime","age","user_score","logincount","buycount","pay","is_sale","spu_pro_status","brand","cate","nprice","noriginal","nsku_num","pro_value"))
.setOutputCol("orign_feature")
val ofdf = va.transform(num_fmt).select($"cust_id",$"good_id",$"flag".alias("label"),$"orign_feature")
// データの正規化処理
val mmScaler = new MinMaxScaler().setInputCol("orign_feature").setOutputCol("features")
val res = mmScaler.fit(ofdf).transform(ofdf)
.select($"cust_id", $"good_id", $"label", $"features")
3、データを2つに分類します:ラベル=0/1を使用して予測、ラベル=2の中の普通ユーザーを使用して推奨
(res.filter("label!=2"),res.filter("label=2"))
2) LRロジスティック回帰:
1、新規ユーザーを抽出し、グローバルホットセールの推奨を行う
アイデア:
普通ID left join 新規+普通ID => new (cust_id,good_id,rank)
val allHot = HDFSConnection.readDataToHDFS(spark,"/myshops/dwd_hotsell")
// 行動のあるユーザーを読み出す
val txt = spark.sparkContext.textFile("file:///D:/logs/virtualLogs/*.log").cache()
import spark.implicits._
val normalUser = txt.map(line=>{
val arr = line.split(" ")
(arr(2),1)
}).toDF("cust_id","flag")
.distinct().cache()
// 新規ユーザーflag=null、新規ユーザーを抽出し、各ユーザーのホットセールのトップ10を求める
// .select($"cust_id",$"good_id",$"rank")
// 全てのユーザー-普通ユーザーflag=1 => 能マッチするflag=1、できないマッチするflag=null
val win = Window.partitionBy("cust_id").orderBy(desc("sellnum"))
// new (cust_id,good_id,rank)
val newUsers = allHot.join(normalUser,Seq("cust_id"),"left")
.filter("flag is null")
// newの前10
val newUserRecommend = newUsers.select($"cust_id",$"good_id",
row_number().over(win).alias("rank"))
.filter(s"rank<=${rank}")
2、predictから新規ユーザーを除外する
アイデア:
($"cust_id", $"good_id", $"label", $"features") left join new => (cust_id,good_id,rank,flag=1) .filter flag = null
val (train,predict):Tuple2[DataFrame,DataFrame] = LRDataHandler.LRdata(spark)
// predictから新規ユーザーを除外する
val newUserID = newUsers.map(x=>{case(cust_id,good_id,rank)=>(cust_id,1)})
.toDF("cust_id","flag")
val normalPredict = predict.join(newUserID,Seq("cust_id"),"left").filter("flag is null")
3、既存のLRモデルが存在しない場合は作成する
// ユーザーが見ている商品を基にLRモデルの訓練データとして使用する
// 首先、HDFS上で既存のLRモデルがあるかどうか確認し、あればそれを取得する
val path = new Path(HDFSConnection.paramMap("hadoop_url")+"/myshops/LR_model"); // HDFS上のファイルパスを指定
val hadoopConf = spark.sparkContext.hadoopConfiguration
val hdfs = org.apache.hadoop.fs.FileSystem.get(new URI(HDFSConnection.paramMap("hadoop_url")+"/myshops/LR_model"),hadoopConf)
var model:LogisticRegressionModel = null;
if(hdfs.exists(path)) {
model = HDFSConnection.readLRModelToHDFS("/myshops/LR_model")
}else{
val lr = new LogisticRegression().setMaxIter(20).setRegParam(0.01)
model = lr.fit(train)
HDFSConnection.writeLRModelToHDFS(model,"/myshops/LR_model")
}
4、LRモデルを使用して、普通ユーザーが見たことのない商品を予測する
// 特徴列を読み取る
// predict中の featuresColのみ
// predictから冷ユーザーを除外する
val res = model.transform(normalPredict).drop("features")
import spark.implicits._
val wnd = Window.partitionBy($"cust_id").orderBy(desc("score"))
// 普通ユーザー recomend 商品 (三路再現後建立 LRモデル) (cust_id,good_id,rank)
val normalRecommend = res.select("cust_id","good_id","probability")
.rdd.map{case(Row(uid:Double,gid:Double,score:DenseVector))=>(uid,gid,score(1))}
.toDF("cust_id","good_id","score")
.select($"cust_id",$"good_id",row_number().over(wnd).alias("rank"))
.filter(s"rank<=${rank}")
5、普通ユーザーと新規ユーザーのLR予測を組み合わせる
// 普通ユーザーと新規ユーザーのrecomend結果を結合 (cust_id,good_id,rank)
val recommend = normalRecommend.union(newUserRecommend)
MYSQLConnection.writeTable(spark,recommend,"userrecommend")