Forráskód Böngészése

增加新的特征表制作

xueyiming 1 napja
szülő
commit
943764993f

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

@@ -48,7 +48,7 @@ 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_level_base_data")
+    val midLevelTable = param.getOrElse("midLevelTable", "alg_recsys_mid_ad_person")
 
     val loader = getClass.getClassLoader
     val resourceUrlBucket = loader.getResource(bucketFile)
@@ -120,7 +120,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
     // 3 循环执行数据生产
     val dateRange = MyDateUtils.getDateRange(beginStr, endStr)
     for (dt <- dateRange) {
-      // 取前一天 mid 分层表,解析 basic_l6.level 作为 user_layer_l6
+      // 取前一天 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,
@@ -130,23 +131,8 @@ object makedata_ad_33_bucketDataFromOriginToHive_20260808 {
           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
-            }
-          }
+          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")