基于Spark/Scala实现联赛球队最长连续未对阵指定队的场次统计
Got it, let's break down how to solve this problem. We need to find, for every team in your football league, the longest streak of consecutive games (both past and future) where they didn't play against MCU. Here's a step-by-step solution using both Spark SQL and Scala DataFrame API, tailored to your table structure.
First: Understand the Data Transformation
Your table has teamname, lastgame (recent opponent), nextgame (upcoming opponents, likely delimited like MCU-BHA-LIV), dateoflastgame, and dateofnextgame. First, we need to flatten this into a single row per game, so we can sort games chronologically and analyze streaks.
Step 1: Create a Unified View of All Games
We'll combine past games (from lastgame) and future games (split from nextgame), making sure each opponent maps to its correct date.
Spark SQL Approach
-- Create a view of all past and future games, one row per game CREATE OR REPLACE VIEW all_games AS SELECT teamname, opponent, game_date FROM ( -- Past games (single row per team) SELECT teamname, lastgame AS opponent, dateoflastgame AS game_date FROM football_table WHERE lastgame IS NOT NULL UNION ALL -- Future games: split delimited strings into individual rows SELECT ft.teamname, opp AS opponent, date_str AS game_date FROM football_table ft -- Split nextgame into individual opponents with positions LATERAL VIEW posexplode(split(ft.nextgame, '-')) AS idx, opp -- Split dateofnextgame into individual dates with positions LATERAL VIEW posexplode(split(ft.dateofnextgame, '-')) AS idx2, date_str -- Match opponent and date by their position in the delimited string WHERE idx = idx2 AND ft.nextgame IS NOT NULL ) ORDER BY teamname, game_date;
Scala DataFrame Approach
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // Load your original table val rawDF = spark.read.table("football_table") // Process past games val pastGamesDF = rawDF .select( col("teamname"), col("lastgame").alias("opponent"), col("dateoflastgame").alias("game_date") ) .filter(col("lastgame").isNotNull) // Process future games: split delimited opponents and dates, then align them val futureGamesDF = rawDF .select( col("teamname"), posexplode(split(col("nextgame"), "-")).alias("idx", "opponent"), posexplode(split(col("dateofnextgame"), "-")).alias("idx2", "game_date") ) .filter(col("idx") === col("idx2") && col("nextgame").isNotNull) .select("teamname", "opponent", "game_date") // Combine past and future games, sorted by team and date val allGamesDF = pastGamesDF.union(futureGamesDF) .orderBy(col("teamname"), col("game_date"))
Step 2: Identify Streaks Using Window Functions
The key idea is to group consecutive non-MCU games together. We'll use a cumulative sum to assign a unique "group ID" to each streak: every time a team plays MCU, we increment the group ID, so all games between two MCU matches (or start/end of schedule) get the same group ID.
Spark SQL Approach
-- Assign group IDs to consecutive non-MCU game streaks CREATE OR REPLACE VIEW game_streak_groups AS SELECT teamname, opponent, game_date, -- Increment group ID every time we encounter an MCU match SUM(CASE WHEN opponent = 'MCU' THEN 1 ELSE 0 END) OVER ( PARTITION BY teamname ORDER BY game_date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS group_id FROM all_games; -- Calculate the longest streak per team SELECT teamname, MAX(CASE WHEN opponent != 'MCU' THEN streak_length ELSE 0 END) AS longest_consecutive_no_mcu FROM ( -- Count the number of games in each group SELECT teamname, group_id, COUNT(*) AS streak_length FROM game_streak_groups GROUP BY teamname, group_id ) GROUP BY teamname ORDER BY longest_consecutive_no_mcu DESC;
Scala DataFrame Approach
// Define a window partitioned by team, ordered by game date val teamWindow = Window.partitionBy("teamname").orderBy("game_date") // Assign group IDs to each streak val groupedGamesDF = allGamesDF .withColumn( "group_id", sum(when(col("opponent") === "MCU", 1).otherwise(0)).over(teamWindow) ) // Calculate streak lengths and find the maximum per team val longestStreaksDF = groupedGamesDF .groupBy("teamname", "group_id") .agg(count("*").alias("streak_length")) .groupBy("teamname") .agg( max(when(col("opponent") =!= "MCU", col("streak_length")).otherwise(0)) .alias("longest_consecutive_no_mcu") ) .orderBy(col("longest_consecutive_no_mcu").desc) // Show the final results longestStreaksDF.show()
Handling Edge Cases
- Teams that never played MCU: The query will return their total number of games as the longest streak (since all games are in one group with no MCU matches).
- Teams that only play MCU: The result will be 0, as there are no non-MCU games.
- Teams with ongoing streaks (end of schedule is non-MCU): The current streak is included in the maximum, which is exactly what we want.
- Teams with no past/future games: The filters (
WHERE lastgame IS NOT NULL/WHERE nextgame IS NOT NULL) ensure we only process existing games, so these teams won't appear in results (or you can adjust to include them with a streak of 0 if needed).
内容的提问来源于stack exchange,提问作者Gerard

