Преглед изворни кода

增加新的特征表制作

xueyiming пре 1 дан
родитељ
комит
f1a64809d8

Разлика између датотеке није приказан због своје велике величине
+ 4 - 0
src/main/resources/20260807_ad_bucket_1112.txt


+ 74 - 63
src/main/scala/com/aliyun/odps/spark/examples/makedata_ad/v20240718/makedata_ad_33_bucketDataFromOriginToHive_20260807.scala → src/main/scala/com/aliyun/odps/spark/examples/makedata_ad/v20240718/makedata_ad_33_bucketDataFromOriginToHive_20260808.scala

@@ -15,7 +15,7 @@ import scala.io.Source
 import scala.language.postfixOps
 import scala.util.Random
 
-object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
+object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
   val CTR_SMOOTH_BETA_FACTOR = 25
   val CVR_SMOOTH_BETA_FACTOR = 10
   val CTCVR_SMOOTH_BETA_FACTOR = 100
@@ -41,13 +41,14 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
     val idDefaultValue = param.getOrElse("idDefaultValue", "1.0").toDouble
     val filterNames = param.getOrElse("filterNames", "").split(",").filter(_.nonEmpty).toSet
     val filterAdverIds = param.getOrElse("filterAdverIds", "").split(",").filter(_.nonEmpty).toSet
-//    val whatLabel = param.getOrElse("whatLabel", "ad_is_conversion")
+    val whatLabel = param.getOrElse("whatLabel", "ad_is_conversion")
     val negSampleRate = param.getOrElse("negSampleRate", "1").toDouble
     // 分割样本集的比例,splitRate部分输出至outputTable,补集输出至outputTable2(如果outputTable2不为空)
     val splitRate = param.getOrElse("splitRate", "0.9").toDouble
     val maskFeatureRate = param.getOrElse("maskFeatureRate", "0.0").toDouble
     val bucketFile = param.getOrElse("bucketFile", "20260807_ad_bucket_1112.txt")
     val flag = param.getOrElse("flag", "0")
+    val midLevelTable = param.getOrElse("midLevelTable", "alg_recsys_mid_ad_level_base_data")
 
     val loader = getClass.getClassLoader
     val resourceUrlBucket = loader.getResource(bucketFile)
@@ -88,17 +89,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
       "user_skuid_click_3d", "user_skuid_click_7d", "user_skuid_click_30d",
       "user_skuid_conver_3d", "user_skuid_conver_7d", "user_skuid_conver_30d",
       "is_weekday", "day_of_the_week", "user_conver_ad_class", "category_name",
-      "material_md5", "user_layer", "user_class", "user_click_ad_class", "user_view_ad_class",
-      "customer", "landing", "flag")
-    val labelFields = Set(
-      "has_click",
-      "has_conversion",
-      "has_scan",
-      "has_addwechat",
-      "has_scan_addwechat",
-      "is_landing3",
-      "is_landing3_with_has_scan"
-    )
+      "material_md5", "user_layer", "user_layer_l6", "user_class", "user_click_ad_class", "user_view_ad_class",
+      "customer", "customer_id", "landing", "landing_page_type", "agent_id", "flag")
 
 
     // 2 读取odps+表信息
@@ -108,16 +100,15 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
 
     // 检查所有字段,收集非法字段
     val invalidFields = tableSchema.flatMap { case (fieldName, _) =>
-      // 如果是标签字段,直接跳过不校验
-      if (labelFields.contains(fieldName)) {
-        None
-      } else {
-        // 否则,校验是否在特征列表里
+      // 跳过 has_click 和 has_conversion 列
+      if (fieldName != "has_click" && fieldName != "has_conversion") {
         if (!lowerCaseDenseFeatureNames.contains(fieldName) && !sparseFeatureNames.contains(fieldName)) {
-          Some(fieldName) // 收集未知字段
+          Some(fieldName) // 收集缺少字段
         } else {
           None
         }
+      } else {
+        None
       }
     }.toList
 
@@ -129,6 +120,36 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
     // 3 循环执行数据生产
     val dateRange = MyDateUtils.getDateRange(beginStr, endStr)
     for (dt <- dateRange) {
+      // 取前一天 mid 分层表,解析 basic_l6.level 作为 user_layer_l6
+      val midLayerDt = MyDateUtils.getNumDaysBefore(dt, 1)
+      val midAdPersonRdd = odpsOps.readTable(
+          project = project,
+          table = midLevelTable,
+          partition = s"dt=$midLayerDt",
+          transfer = func,
+          numPartition = tablePart
+        ).flatMap(record => {
+          val midVal = Option(record.getString("mid")).getOrElse("")
+          val featureStr = Option(record.getString("feature")).getOrElse("")
+          if (midVal.isEmpty || featureStr.isEmpty) {
+            None
+          } else {
+            try {
+              val featureJson = JSON.parseObject(featureStr)
+              val basicL6 =
+                if (featureJson == null || !featureJson.containsKey("basic_l6")) null
+                else featureJson.getJSONObject("basic_l6")
+              val level =
+                if (basicL6 == null || !basicL6.containsKey("level")) null
+                else basicL6.getString("level")
+              if (level != null && level.nonEmpty) Some((midVal, level)) else None
+            } catch {
+              case _: Exception => None
+            }
+          }
+        }).reduceByKey((a, _) => a).cache()
+      println(s"load midLevelTable=$midLevelTable dt=$midLayerDt for sampleDt=$dt")
+
       val timeRange = MyDateUtils.getDateHourRange(dt + "06", dt + "23")
       val recordRdd = timeRange.map { dt_hh =>
           val dt = dt_hh.substring(0, 8)
@@ -166,17 +187,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
               !filterAdverIds.contains(adverId)
             })
             .filter(record => {
-              val labelJson = getJsonObject(record, "label_json")
-              val extendAlg = getJsonObject(record, "extend_alg")
-              val hasScanLabel = labelJson.getString("ad_is_scan").toInt
-              val hasConversionLabel = labelJson.getString("ad_is_conversion").toInt
-              val isLanding3 = extendAlg.getString("landing_page_type") == "3"
-              val randVal = Random.nextDouble()
-              if (isLanding3) {
-                (hasScanLabel > 0) || (randVal < negSampleRate)
-              } else {
-                hasConversionLabel > 0 || (randVal < negSampleRate)
-              }
+              val label = record.getString(whatLabel).toInt
+              label > 0 || Random.nextDouble() < negSampleRate
             })
             .map(record => {
               val featureMap = new JSONObject()
@@ -193,7 +205,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
               val mid = record.getString("mid")
               val pqtid = record.getString("pqtid")
               val apptype = record.getString("apptype")
-              val targetingConversion = record.getString("targeting_conversion")
+              val targetingConversion = Option(record.getString("targeting_conversion")).getOrElse("")
+
               featureMap.put("apptype", apptype)
               featureMap.put("ts", ts)
               featureMap.put("mid", mid)
@@ -264,8 +277,9 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
               if (b1.containsKey("creative_action_embedding") && b1.getString("creative_action_embedding").nonEmpty) {
                 featureMap.put("creative_action_embedding", b1.getString("creative_action_embedding").split('|').map(_.toDouble).map(_.toFloat).mkString("|"))
               }
-              if (extendAlg.containsKey("customer_id")) {
+              if (extendAlg.containsKey("customer_id") && extendAlg.getString("customer_id").nonEmpty) {
                 featureMap.put("customer", extendAlg.getString("customer_id"))
+                featureMap.put("customer_id", extendAlg.getString("customer_id"))
               }
               if (sceneFeature.containsKey("hour") && sceneFeature.getString("hour").nonEmpty) {
                 featureMap.put("hour", sceneFeature.getString("hour"))
@@ -304,8 +318,12 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
                   featureMap.put(key, value)
                 }
               }
-              if (extendAlg.containsKey("landing_page_type")) {
+              if (extendAlg.containsKey("landing_page_type") && extendAlg.getString("landing_page_type").nonEmpty) {
                 featureMap.put("landing", extendAlg.getString("landing_page_type"))
+                featureMap.put("landing_page_type", extendAlg.getString("landing_page_type"))
+              }
+              if (reqFeature.containsKey("agentId") && reqFeature.getString("agentId").nonEmpty) {
+                featureMap.put("agent_id", reqFeature.getString("agentId"))
               }
               if (reqFeature.containsKey("layer_l4")) {
                 featureMap.put("user_layer", reqFeature.getString("layer_l4"))
@@ -536,7 +554,7 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
                 featureMap.put("user_vid_return_tags_14d", e1.getString("tags_14d"))
               }
 
-              if (e2.containsKey("tags_1d") && e2.getString("tags_1d").nonEmpty) {
+              if (e2.containsKey("tags_14d") && e2.getString("tags_14d").nonEmpty) {
                 featureMap.put("user_vid_share_tags_1d", e2.getString("tags_1d"))
               }
               if (e2.containsKey("tags_14d") && e2.getString("tags_14d").nonEmpty) {
@@ -694,15 +712,14 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
                 ("k1", k1), ("k2", k2), ("k3", k3), ("k4", k4)
               )
               val kPeriods = List("2h", "4h", "6h", "12h", "1d", "3d", "today", "1w")
-              val eventId = Option(targetingConversion).getOrElse("")
               for ((kPrefix, kn) <- kList) {
                 for (period <- kPeriods) {
                   val view = if (kn.isEmpty) 0D else kn.getIntValue("ad_view_" + period).toDouble
                   val click = if (kn.isEmpty) 0D else kn.getIntValue("ad_click_" + period).toDouble
                   val eventJson = parseEventJson(kn, "event_" + period)
                   val conver =
-                    if (eventJson.isEmpty || eventId.isEmpty) 0D
-                    else eventJson.getIntValue(eventId + "_" + period).toDouble
+                    if (eventJson.isEmpty || targetingConversion.isEmpty) 0D
+                    else eventJson.getIntValue(targetingConversion + "_" + period).toDouble
                   val ctr = RankExtractorFeature_20240530.divSmooth2(click, view, CTR_SMOOTH_BETA_FACTOR)
                   val cvr = RankExtractorFeature_20240530.divSmooth2(conver, click, CVR_SMOOTH_BETA_FACTOR)
                   val ctvr = RankExtractorFeature_20240530.divSmooth2(conver, view, CTCVR_SMOOTH_BETA_FACTOR)
@@ -753,25 +770,10 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
 
               //4 处理label信息。
               val labels = new JSONObject
-              val landing = featureMap.getString("landing")
-              val labelJson: JSONObject = getJsonObject(record, "label_json") // label_json
-              if (!labelJson.isEmpty) {
-                labels.put("has_scan", labelJson.getInteger("ad_is_scan"))
-                labels.put("has_addwechat", labelJson.getInteger("ad_is_addwechat"))
-                labels.put("has_scan_addwechat", labelJson.getInteger("ad_is_addwechat"))
-                if (landing == "3") {
-                  labels.put("is_landing3", 1)
-                  if (labelJson.getIntValue("ad_is_scan") == 1 && targetingConversion == "10004") {
-                    labels.put("is_landing3_with_has_scan", 1)
-                  } else {
-                    labels.put("is_landing3_with_has_scan", 0)
-                  }
-                } else {
-                  labels.put("is_landing3", 0)
-                  labels.put("is_landing3_with_has_scan", 0)
+              for (labelKey <- List("ad_is_click", "ad_is_conversion")) {
+                if (!record.isNull(labelKey)) {
+                  labels.put(labelKey, record.getString(labelKey))
                 }
-                labels.put("has_click", labelJson.getIntValue("ad_is_click"))
-                labels.put("has_conversion", labelJson.getIntValue("ad_is_conversion"))
               }
               //5 处理log key表头。
               val headvideoid = record.getString("headvideoid")
@@ -781,6 +783,19 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
             })
           odpsData
         }).reduce(_ union _)
+        .map { case (logKey, labelKey, featureMap) =>
+          val midKey = Option(featureMap.getString("mid")).getOrElse("")
+          (midKey, (logKey, labelKey, featureMap))
+        }
+        .leftOuterJoin(midAdPersonRdd)
+        .map {
+          case (_, ((logKey, labelKey, featureMap), Some(adPerson))) if adPerson != null && adPerson.nonEmpty =>
+            featureMap.put("user_layer_l6", adPerson)
+            (logKey, labelKey, featureMap)
+          case (_, ((logKey, labelKey, featureMap), _)) =>
+            featureMap.put("user_layer_l6", "无曝光")
+            (logKey, labelKey, featureMap)
+        }
         .map { case (logKey, labelKey, jsons) =>
           val denseFeatures = scala.collection.mutable.Map[String, Double]()
           val sparseFeatures = scala.collection.mutable.Map[String, String]()
@@ -799,7 +814,7 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
         .map {
           case (logKey, labelKey, denseFeatures, sparseFeatures) =>
             val labelObject = JSON.parseObject(labelKey)
-//            val label = labelObject.getOrDefault(whatLabel, "0").toString
+            val label = labelObject.getOrDefault(whatLabel, "0").toString
             val bucketsMap = bucketsMap_br.value
             var resultMap = denseFeatures.collect {
               case (name, score) if !filterNames.exists(name.contains) && score > 1E-8 =>
@@ -814,13 +829,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
             sparseFeatures.foreach(kv => {
               resultMap += (kv._1 -> kv._2)
             })
-            resultMap += ("has_click" -> labelObject.getString("has_click"))
-            resultMap += ("has_conversion" -> labelObject.getString("has_conversion"))
-            resultMap += ("has_scan" -> labelObject.getString("has_scan"))
-            resultMap += ("has_addwechat" -> labelObject.getString("has_addwechat"))
-            resultMap += ("has_scan_addwechat" -> labelObject.getString("has_scan_addwechat"))
-            resultMap += ("is_landing3" -> labelObject.getString("is_landing3"))
-            resultMap += ("is_landing3_with_has_scan" -> labelObject.getString("is_landing3_with_has_scan"))
+            resultMap += ("has_click" -> labelObject.getString("ad_is_click"))
+            resultMap += ("has_conversion" -> labelObject.getString("ad_is_conversion"))
             resultMap += ("logkey" -> logKey)
             resultMap
         }.coalesce(128)
@@ -834,6 +844,7 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260807 {
         odpsOps.saveToTable(project, outputTable, partition, splitRdds(0), write, defaultCreate = true, overwrite = true)
         odpsOps.saveToTable(project, outputTable2, partition, splitRdds(1), write, defaultCreate = true, overwrite = true)
       }
+      midAdPersonRdd.unpersist()
     }
   }
 

Неке датотеке нису приказане због велике количине промена