code
import org.apache.spark.sql.functions._val scored = df
.withColumn("brand_norm", lower(trim(col("brand"))))
.groupBy("day","engine","brand_norm")
code
.agg(countDistinct("query_id").as("queries")
)
.withColumn("mentions_per_query", col("mentions") / col("queries"))
scored.write.mode("overwrite").parquet("s3://企业AI落地/agg/visibility_v3")
Spark作业贵,所以我尽量增量:只重算受影响品牌。企业AI落地指标口径一变,老板第一句就是「历史能对比吗」,所以版本化很重要。
你们有没有用dbt?我想把指标层再规范一点。
(场景参考:广州本地企业试点)