Browse Source

修改人群分层获取方式

xueyiming 1 ngày trước cách đây
mục cha
commit
9b8e5d6a08

+ 1 - 1
src/main/resources/20260807_ad_feature_name.txt

@@ -189,4 +189,4 @@ k4_1w_click
 k4_1w_conver
 k4_1w_ctr
 k4_1w_cvr
-k4_1w_ctvr
+k4_1w_ctvr

+ 5 - 31
src/main/scala/com/aliyun/odps/spark/examples/makedata_ad/v20240718/makedata_ad_33_bucketDataFromOriginToHive_20260808.scala

@@ -48,7 +48,6 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
     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_person")
 
     val loader = getClass.getClassLoader
     val resourceUrlBucket = loader.getResource(bucketFile)
@@ -120,22 +119,6 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
     // 3 循环执行数据生产
     val dateRange = MyDateUtils.getDateRange(beginStr, endStr)
     for (dt <- dateRange) {
-      // 取前一天 mid 分层表(STRING 中间表,避免源表 feature JSON 类型 Tunnel 读失败)
-      // 表结构:mid STRING, ad_person STRING, PARTITIONED BY (dt)
-      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 adPerson = Option(record.getString("ad_person")).getOrElse("")
-          if (midVal.nonEmpty && adPerson.nonEmpty) Some((midVal, adPerson)) else 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)
@@ -315,6 +298,11 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
               if (reqFeature.containsKey("layer_l4")) {
                 featureMap.put("user_layer", reqFeature.getString("layer_l4"))
               }
+              if (reqFeature.containsKey("layer_l6") && reqFeature.getString("layer_l6").nonEmpty) {
+                featureMap.put("user_layer_l6", reqFeature.getString("layer_l6"))
+              } else {
+                featureMap.put("user_layer_l6", "无曝光")
+              }
               if (reqFeature.containsKey("clazz_l4")) {
                 featureMap.put("user_class", reqFeature.getString("clazz_l4"))
               }
@@ -773,19 +761,6 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
             })
           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]()
@@ -838,7 +813,6 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
         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()
     }
   }