{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":31254,"databundleVersionId":3103714,"sourceType":"competition"}],"dockerImageVersionId":30587,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os\n!pip install pyspark\nfrom pyspark.sql.session import SparkSession\nsc = SparkSession.builder.appName(\"Recommendations\").config(\"spark.sql.files.maxPartitionBytes\", 500000000).getOrCreate()\nspark = SparkSession(sc)","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2023-11-20T06:09:04.830717Z","iopub.execute_input":"2023-11-20T06:09:04.831116Z","iopub.status.idle":"2023-11-20T06:10:01.637012Z","shell.execute_reply.started":"2023-11-20T06:09:04.831084Z","shell.execute_reply":"2023-11-20T06:10:01.635027Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import pyspark\nimport numpy as np\nimport pandas as pd\nfrom pyspark.sql.window import *\nfrom pyspark.sql.window import Window\nfrom pyspark.sql.functions import *\nfrom pyspark.sql.functions import col, row_number, udf,when\nfrom pyspark.sql.functions import to_timestamp,date_format\nfrom pyspark.sql.functions import array_contains\nfrom pyspark.sql.types import *\nfrom pyspark.sql.types import StringType, ArrayType \nfrom pyspark.sql.types import StructType,StructField,IntegerType\nfrom pyspark.sql.types import DoubleType, BooleanType\nfrom pyspark.sql import SQLContext \nfrom pyspark.ml.recommendation import ALS\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.ml.feature import StringIndexer\nfrom pyspark.ml import Pipeline\nfrom pyspark.ml.tuning import CrossValidator, ParamGridBuilder","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"transaction = spark.read.option(\"header\",True).csv(\"../input/h-and-m-personalized-fashion-recommendations/transactions_train.csv\")\ncustomer = spark.read.option(\"header\",True).csv(\"../input/h-and-m-personalized-fashion-recommendations/customers.csv\")\narticles = spark.read.option(\"header\",True).csv(\"../input/h-and-m-personalized-fashion-recommendations/articles.csv\")","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:10:27.150540Z","iopub.execute_input":"2023-11-20T06:10:27.152141Z","iopub.status.idle":"2023-11-20T06:10:34.671620Z","shell.execute_reply.started":"2023-11-20T06:10:27.152073Z","shell.execute_reply":"2023-11-20T06:10:34.670467Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"Step 1 Process the dataset","metadata":{}},{"cell_type":"code","source":"hm = transaction.withColumn('t_dat', transaction['t_dat'].cast('string'))\nhm = hm.withColumn('date', from_unixtime(unix_timestamp('t_dat', 'yyyy-MM-dd')))\nhm = hm.withColumn('year', year(col('date')))\nhm = hm.withColumn('month', month(col('date')))\n#hm = hm.withColumn('day', date_format(col('date'), \"d\"))\nhm = hm.withColumn('day', day(col('date')))\n\nhm1 = hm.filter(\"year == 2020 and month == 9 and day <= 15\")\nhm2 = hm1.select(\"customer_id\",\"article_id\")\nhm3 = hm2.groupBy(\"customer_id\",\"article_id\").count() #.sort(\"count\", ascending = False)\nhm4 = hm3.sort(\"count\", ascending = False)","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:10:40.344388Z","iopub.execute_input":"2023-11-20T06:10:40.344793Z","iopub.status.idle":"2023-11-20T06:10:40.754955Z","shell.execute_reply.started":"2023-11-20T06:10:40.344762Z","shell.execute_reply":"2023-11-20T06:10:40.753657Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"top30_rank = hm4.drop(\"customer_id\").head(30)\ntop30_rank = [(one[0],one[1]) for one in top30_rank]\ntop30 = [one[0] for one in top30_rank]\nitem_type = articles.select('product_type_name', 'article_id')\nprint(top30)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"Step 2 Get type popularity suggestions","metadata":{}},{"cell_type":"code","source":"item_trans = hm3.join(item_type, hm3.article_id == item_type.article_id, \"leftouter\")\nitem_trans = item_trans.drop(item_type.article_id)\nitem_type_count = item_trans.groupBy(\"article_id\",\"product_type_name\").sum(\"count\").withColumnRenamed(\"sum(count)\", \"count\")\nwindow1 = Window.partitionBy(\"product_type_name\").orderBy(col(\"count\").desc())\nwindow1_sorted = item_type_count.select(\"product_type_name\",\"article_id\",\"count\").withColumn(\"row\",row_number().over(window1)) \\\n  .filter(col(\"row\") <= 6) \ntop_type = item_type_count.groupBy(\"product_type_name\").sum(\"count\").withColumnRenamed(\"sum(count)\", \"count\").sort(\"count\", ascending = False)\ntop5type = top_type.head(5)\ntop10type = top_type.head(10)\nuser_type = item_trans.groupBy(\"customer_id\",\"product_type_name\").sum(\"count\").withColumnRenamed(\"sum(count)\",\"count\")\nwindow2 = Window.partitionBy(\"customer_id\").orderBy(col(\"count\").desc())\nwindow2_sorted = user_type.select(\"customer_id\",\"product_type_name\",\"count\").withColumn(\"rank\",row_number().over(window2)) \\\n  .filter(col(\"rank\") <= 3)\ntype_item_lw = window1_sorted.drop(\"count\").withColumn(\"item_and_rank\", struct(window1_sorted.article_id,window1_sorted.row))\ntype_item_lw2 = type_item_lw.groupBy(\"product_type_name\").agg(collect_list(col(\"item_and_rank\")))\nuser_type_lw = window2_sorted.drop(\"count\")\ntype_item_lw3 = type_item_lw2.withColumnRenamed(\"collect_list(item_and_rank)\", \"items_list\")\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def itemlist_to_list(x):\n    listx = x\n    sizex = len(listx)\n    newlist = []\n    for one in listx:\n        article_id = one[\"article_id\"]\n        row = one[\"row\"]\n        newlist.append((article_id,row))\n    # row in ascending order, 1,2,3,4, and so on.\n    newlist.sort(key = lambda x: x[1])\n    finallist = []\n    for one in newlist:\n        finallist.append(one[0])\n    return finallist\n    ","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"itemlistUDF = udf(lambda x:itemlist_to_list(x),ArrayType(StringType())) \ntype_item_lw4 = type_item_lw3.withColumn(\"listitems\",itemlistUDF(col(\"items_list\")))\ntype_item_lw5 = type_item_lw4.drop(\"items_list\")\nuser_type_item = user_type_lw.join(type_item_lw5, user_type_lw.product_type_name == type_item_lw5.product_type_name, \"leftouter\" )\nuser_type_item = user_type_item.drop(type_item_lw5.product_type_name)\nuser_type_item1 = user_type_item.withColumn(\"type_and_item_and_rank\", struct(user_type_item.product_type_name, user_type_item.listitems, user_type_item.rank))\nuser_type_item2 = user_type_item1.groupBy(\"customer_id\").agg(collect_list(col(\"type_and_item_and_rank\")))\nuser_type_item3 = user_type_item2.withColumnRenamed(\"collect_list(type_and_item_and_rank)\", \"list_type_item_rank\")\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def extend_listitems(x):\n    if x == None:\n        return []\n    if len(x) == 0:\n        return []\n    listx = x\n    newlist = []\n    for one in listx:\n        listitems = one['listitems']\n        rank = one[\"rank\"]\n        newlist.append((listitems,rank))\n    newlist.sort(key = lambda x: x[1])\n    finallist = []\n    for one in newlist:\n        finallist.extend(one[0])\n    return finallist","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"typeitemlistUDF = udf(lambda x:extend_listitems(x),ArrayType(StringType())) \nuser_type_item4 = user_type_item3.withColumn(\"list_items\",typeitemlistUDF(col(\"list_type_item_rank\")))\nuser_type_item4 = user_type_item4.drop(\"list_type_item_rank\")\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"Step 3 Buy Agains","metadata":{}},{"cell_type":"code","source":"hm5 = hm4.filter(\"count > 1\")\nwindow3 = Window.partitionBy(\"customer_id\").orderBy(col(\"count\").desc())\nwindow3_sorted = hm5.withColumn(\"rank\",row_number().over(window3)) \\\n  .filter(col(\"rank\") <= 5)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rebuy1 = window3_sorted.filter(\"rank <= 5\").drop(\"count\")\nrebuy2 = rebuy1.withColumn(\"item_rank\",struct(rebuy1.article_id,rebuy1.rank))\nrebuy3 = rebuy2.groupBy(\"customer_id\").agg(collect_list(col(\"item_rank\")))\nrebuy4 = rebuy3.withColumnRenamed(\"collect_list(item_rank)\", \"item_list_rank\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def extend_list_items(x):\n    if x == None:\n        return []\n    if len(x) == 0:\n        return []\n    listx = x\n    newlist = []\n    for one in listx:\n        item = one['article_id']\n        rank = one[\"rank\"]\n        newlist.append((item,rank))\n    newlist.sort(key = lambda x: x[1])\n    finallist = []\n    for one in newlist:\n        finallist.append(one[0])\n    return finallist","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"typeitem_listUDF = udf(lambda x:extend_list_items(x),ArrayType(StringType())) \nrebuy5 = rebuy4.withColumn(\"list_items\",typeitem_listUDF(col(\"item_list_rank\")))\nrebuy6 = rebuy5.drop(\"item_list_rank\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rebuy_and_pop_type = rebuy6.union(user_type_item4)\nrebuy_and_pop_type2 = rebuy_and_pop_type.groupBy(\"customer_id\").agg(collect_list(col(\"list_items\")))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def extend_list(x):\n    if x == None:\n        return []\n    if len(x) == 0:\n        return []\n    listx = x\n    newlist = []\n    for one in listx:\n        newlist.extend(one)\n    return newlist","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"combine_listUDF = udf(lambda x:extend_list(x),StringType()) ","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"\nrebuy_and_pop_type3 = rebuy_and_pop_type2.withColumn(\"all_list_items\",combine_listUDF(col(\"collect_list(list_items)\")))\nrebuy_and_pop_type3 = rebuy_and_pop_type3.drop(\"collect_list(list_items)\")\nrebuy_pop = rebuy_and_pop_type3.select(concat_ws(\" : \", rebuy_and_pop_type3.customer_id, rebuy_and_pop_type3.all_list_items)\n              .alias(\"all_output\"))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rebuy_pop.write.text(\"/kaggle/working/rebuy_pop2.txt\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"Step 4 Model-based Collaborative Filtering","metadata":{}},{"cell_type":"code","source":"hm1 = hm.filter(\"year == 2020 and month == 9 and day <= 15\")\nhm_2020_09 = hm1.filter(\"day >= 13\") # day >= 9\n# Prepare the dataset to get customer buy how many articles\nhm_df = hm_2020_09.groupby('customer_id', 'article_id').count()","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:11:09.600500Z","iopub.execute_input":"2023-11-20T06:11:09.600905Z","iopub.status.idle":"2023-11-20T06:11:09.661253Z","shell.execute_reply.started":"2023-11-20T06:11:09.600872Z","shell.execute_reply":"2023-11-20T06:11:09.660063Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"indexer = [StringIndexer(inputCol=column, outputCol=column+\"_index\") for column in list(set(hm_df.columns)-set(['count'])) ]\npipeline = Pipeline(stages=indexer)\ntransformed = pipeline.fit(hm_df).transform(hm_df)\n#transformed.show()\n(training,test)=transformed.randomSplit([0.8, 0.2])","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:11:17.278881Z","iopub.execute_input":"2023-11-20T06:11:17.279276Z","iopub.status.idle":"2023-11-20T06:13:33.961164Z","shell.execute_reply.started":"2023-11-20T06:11:17.279242Z","shell.execute_reply":"2023-11-20T06:13:33.960170Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# get unique customer and article ids to map to int\n\"\"\"\nunique_customers = transaction.select('customer_id').distinct()\nunique_articles = transaction.select('article_id').distinct()\ncustomer_id_mapping = []\nfor i, c in enumerate(unique_customers.collect()):\n    customer_id_mapping.append(Row(c['customer_id'], i))\narticle_id_mapping = []\nfor i, c in enumerate(unique_articles.collect()):\n    article_id_mapping.append(Row(c['article_id'], i))\ncustomer_map = spark.sparkContext.createDataFrame(customer_id_mapping, ['customer_id', 'int_customer_id'])\narticle_map = spark.sparkContext.createDataFrame(article_id_mapping, ['article_id', 'int_article_id'])\nmap_df = transaction.join(customer_map, 'customer_id').join(article_map, on='article_id')\nuser_item_ml = map_df.groupby(['int_customer_id', 'int_article_id']).count()\n\n(training,test)=user_item_ml.randomSplit([0.8, 0.2])\n\"\"\"","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"als=ALS(userCol=\"customer_id_index\",itemCol=\"article_id_index\",ratingCol=\"count\",coldStartStrategy=\"drop\",nonnegative=True,implicitPrefs = True)\nals_simple=ALS(maxIter= 15,regParam=0.09,rank=15,userCol=\"customer_id_index\",itemCol=\"article_id_index\",ratingCol=\"count\",coldStartStrategy=\"drop\",nonnegative=True,implicitPrefs = True)","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:14:02.978562Z","iopub.execute_input":"2023-11-20T06:14:02.978956Z","iopub.status.idle":"2023-11-20T06:14:03.014948Z","shell.execute_reply.started":"2023-11-20T06:14:02.978916Z","shell.execute_reply":"2023-11-20T06:14:03.014081Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"param_grid = ParamGridBuilder()\\\n            .addGrid(als.rank, [15,20])\\\n            .addGrid(als.maxIter,[15,20])\\\n            .addGrid(als.regParam,[0.05,0.09,0.14])\\\n            .build()","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:14:06.586768Z","iopub.execute_input":"2023-11-20T06:14:06.587189Z","iopub.status.idle":"2023-11-20T06:14:06.593208Z","shell.execute_reply.started":"2023-11-20T06:14:06.587156Z","shell.execute_reply":"2023-11-20T06:14:06.591830Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"evaluator = RegressionEvaluator(metricName=\"rmse\", labelCol=\"count\",\n                                predictionCol=\"prediction\")","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:14:09.866894Z","iopub.execute_input":"2023-11-20T06:14:09.867289Z","iopub.status.idle":"2023-11-20T06:14:09.882047Z","shell.execute_reply.started":"2023-11-20T06:14:09.867259Z","shell.execute_reply":"2023-11-20T06:14:09.880598Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#als=ALS(userCol=\"customer_id_index\",itemCol=\"article_id_index\",ratingCol=\"count\",coldStartStrategy=\"drop\",nonnegative=True,implicitPrefs = True)\n#als_simple=ALS(maxIter= 15,regParam=0.09,rank=25,userCol=\"customer_id_index\",itemCol=\"article_id_index\",ratingCol=\"count\",coldStartStrategy=\"drop\",nonnegative=True,implicitPrefs = True)\ncv = CrossValidator(estimator=als,estimatorParamMaps=param_grid, evaluator=evaluator,numFolds=3)","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:14:21.501617Z","iopub.execute_input":"2023-11-20T06:14:21.501997Z","iopub.status.idle":"2023-11-20T06:14:21.507630Z","shell.execute_reply.started":"2023-11-20T06:14:21.501969Z","shell.execute_reply":"2023-11-20T06:14:21.506466Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"model = cv.fit(training)","metadata":{"execution":{"iopub.status.busy":"2023-11-20T06:14:27.581129Z","iopub.execute_input":"2023-11-20T06:14:27.581528Z","iopub.status.idle":"2023-11-20T06:18:16.580792Z","shell.execute_reply.started":"2023-11-20T06:14:27.581495Z","shell.execute_reply":"2023-11-20T06:18:16.578835Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"model_2020_09=als_simple.fit(training)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Extract best model from the tuning exercise using ParamGridBuilder\nbest_model = model.bestModel\n\n#Generate predictions and evaluate using RMSE\npredictions = best_model.transform(test)\nrmse = evaluator.evaluate(predictions)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#print evaluation metrics and model parameters\nprint(\"RMSE =\" + str(rmse))\nprint(\"**Best Model**\")\nprint(\"Rank : {}\".format(best_model.rank))\nprint(\"MaxIter: {}\".format(best_model._java_obj.parent().getMaxIter()))\nprint(\"RegParam: {}\".format(best_model._java_obj.parent().getRegParam()))\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"predictions = model_2020_09.transform(test)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"evaluator = RegressionEvaluator(metricName=\"rmse\", labelCol=\"count\",\n                                predictionCol=\"prediction\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rmse = evaluator.evaluate(predictions)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rmse","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_recs_2020_9=model_2020_09.recommendForAllUsers(10)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_recs_2020_9.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.sql.types import IntegerType\ndef id_only_list(x):\n    xlist = x\n    newlist = []\n    for one in xlist:\n        int_id = one[0]\n        newlist.append(int_id)\n    return newlist","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"idonlylistUDF = udf(lambda x:id_only_list(x),ArrayType(IntegerType()))\nuser_ml_rec = user_recs_2020_9.withColumn(\"list_items\",idonlylistUDF(col(\"recommendations\")))\nuser_ml_rec2 = user_ml_rec.drop(\"recommendations\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#user_ml_rec.show(3)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec2.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec2.show(3)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#int_customer_id integer, list of integer: list_items\n# Goal: change these integers to strings\n#user_ml_rec3 = user_ml_rec2.withColumn(\"in_article_id\", col(\"list_items\"))#map(lambda row: (row[\"int_customer_id\"],row[list_items][0],row[list_items][1])).toDF([\"int_customer_id\",\"int_article_id\",\"rank\"])\ndef changing_row(x):\n    user_id = x[0]\n    list_items = x[1]\n    order = 0\n    newlist = []\n    for one in list_items:\n        newlist.append((user_id,one,order))\n        order += 1\n    return newlist\nuser_ml_rec3 = user_ml_rec2.rdd.flatMap(changing_row).toDF([\"int_customer_id\",\"int_article_id\",\"rank\"])\n# function to change to lsit of article_id\n# done.","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec3.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"md=transformed.select(transformed['article_id'],transformed['article_id_index'],transformed['customer_id'],transformed['customer_id_index'])\narticle_map = md.select(\"article_id\",\"article_id_index\").withColumnRenamed('article_id_index', 'int_article_id').distinct()\ncustomer_map = md.select(\"customer_id\",\"customer_id_index\").withColumnRenamed('customer_id_index', 'int_customer_id').distinct()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec4 = user_ml_rec3.join(customer_map, 'int_customer_id').join(article_map, on='int_article_id')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec5 = user_ml_rec4.withColumn(\"item_list_struct\", struct(user_ml_rec4.article_id,user_ml_rec4.rank))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec6 = user_ml_rec5.groupBy(\"customer_id\").agg(collect_list(col(\"item_list_struct\")))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec7 = user_ml_rec6.withColumnRenamed(\"collect_list(item_list_struct)\", \"list_items\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec7.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def ml_itemlist(x):\n    if x == None:\n        return []\n    if len(x) == 0:\n        return []\n    listx = x\n    newlist = []\n    for one in listx:\n        item = one['article_id']\n        rank = one[\"rank\"]\n        newlist.append((item,rank))\n    newlist.sort(key = lambda x: x[1])\n    finallist = []\n    for one in newlist:\n        finallist.append(one[0])\n    return finallist","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec7.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"ml_itemlistUDF = udf(lambda x:ml_itemlist(x),StringType()) \nuser_ml_rec8 = user_ml_rec7.withColumn(\"all_list_items\",ml_itemlistUDF(col(\"list_items\")))\n#user_ml_rec8 = user_ml_rec8.drop(\"collect_list(list_items)\")\nml_recom = user_ml_rec8.select(concat_ws(\" : \", user_ml_rec8.customer_id, user_ml_rec8.all_list_items)\n              .alias(\"all_output\"))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"ml_recom.show(3)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"ml_recom.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"ml_recom.write.text(\"/kaggle/working/ml_recom2.txt\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec7.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"ml_itemlistUDF2 = udf(lambda x:ml_itemlist(x),ArrayType(StringType()))\nuser_ml_rec10 = user_ml_rec7.withColumn(\"all_list_items\",ml_itemlistUDF2(col(\"list_items\")))\nuser_ml_rec10 = user_ml_rec10.drop(\"list_items\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"user_ml_rec10.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"combine_listUDF2 = udf(lambda x:extend_list(x),ArrayType(StringType()))\nrebuy_and_pop_type4 = rebuy_and_pop_type2.withColumn(\"all_list_items\",combine_listUDF2(col(\"collect_list(list_items)\")))\nrebuy_and_pop_type4 = rebuy_and_pop_type4.drop(\"collect_list(list_items)\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"rebuy_and_pop_type4.printSchema()\nuser_ml_rec10.printSchema()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# last join: rebuy, popularity, model\nfinal1 = rebuy_and_pop_type4.union(user_ml_rec10)\nfinal2 = final1.groupBy(\"customer_id\").agg(collect_list(col(\"all_list_items\")))\nfinal3 = final2.withColumn(\"every_list_items\",combine_listUDF(col(\"collect_list(all_list_items)\")))\nfinal3 = final3.drop(\"collect_list(all_list_items)\")\nfinal_out = final3.select(concat_ws(\" : \", final3.customer_id, final3.every_list_items)\n              .alias(\"all_output\"))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"final_out.write.text(\"/kaggle/working/final_out.txt\")","metadata":{"trusted":true},"execution_count":null,"outputs":[]}]}