Spark中withColumn仅返回新增列且计算30分钟浏览量时报错求助
Hey there! Let's tackle your two Spark issues one by one, starting with the simpler one first.
.withColumn() Normally, .withColumn() should add a new column while preserving all existing columns in your DataFrame. If you're only seeing the new column, there's almost certainly a step before or after this call that's trimming your column set. Here are some common checks:
- Did you accidentally use
.select()or.drop()before calling.withColumn()that removed most columns? Double-check the full lineage of your DataFrame leading up to this operation. - Are you reassigning the result incorrectly? For example, if you chain
.select(newColumnName)right after.withColumn(), that would explicitly filter to only the new column. - Have you enabled any unusual Spark configurations that might alter column retention? This is unlikely, but worth ruling out if the above checks don't resolve it.
Let's break down why this error is happening by looking at your code flow:
- You start with
withOneDayts, which contains columns likecid,timestamp,15_min_ts,30_min_ts, etc. - When creating
with15MinViews, you run.groupBy("cid").agg(sum("views"))— groupBy + aggregation collapses your DataFrame to only include the grouped column (cid) and the aggregated result. Then you rename the sum toviews_in_15_minand drop15_min_ts, leavingwith15MinViewswith only two columns:cidandviews_in_15_min. - Next, you try to call
.withColumn("views", getViews(with15MinViews("timestamp"), with15MinViews("30_min_ts")))onwith15MinViews— but this DataFrame no longer hastimestampor30_min_ts! That's exactly why Spark throws the column resolution error.
Fix for this issue
Instead of chaining the 30-minute calculation onto the aggregated with15MinViews, calculate both the 15-minute and 30-minute views separately from the original withOneDayts, then join the results together on cid. Here's adjusted code that works:
// Calculate 15-minute views from the original withOneDayts val with15MinViews = withOneDayts .withColumn("views", getViews(withOneDayts("timestamp"), withOneDayts("15_min_ts"))) .groupBy("cid") .agg(sum("views").alias("views_in_15_min")) // Use alias directly for cleaner code // Calculate 30-minute views separately from withOneDayts val with30MinViews = withOneDayts .withColumn("views", getViews(withOneDayts("timestamp"), withOneDayts("30_min_ts"))) .groupBy("cid") .agg(sum("views").alias("views_in_30_min")) // Join the two aggregated DataFrames to get both metrics in one place val finalResult = with15MinViews.join(with30MinViews, Seq("cid"), "inner")
This approach ensures both calculations use the full withOneDayts DataFrame (which has all required columns like timestamp), then combines the aggregated results via a join. No more missing column errors!
内容的提问来源于stack exchange,提问作者NamrataK

