If you know window functions in T-SQL, you already know most of Spark SQL's window functions. The syntax is nearly identical. The differences are in the defaults, a few missing conveniences, and how the work runs on a cluster. This post goes through the patterns I use most on delegated authority bordereaux, from Spark SQL transformations in AWS Glue and Databricks, with the T-SQL habits that need adjusting.
All examples run as-is in a Fabric or Databricks notebook (%%sql cell) or through spark.sql().
Sample data: a premium bordereau
Each row is a premium transaction reported by a coverholder against a binding authority contract section (identified by UMR and section number). Coverholders resubmit corrected bordereaux, so the same transaction can arrive more than once.
1: CREATE OR REPLACE TEMP VIEW premium_bdx AS 2: SELECT * FROM VALUES 3: ('B0999CH001', 'S1', 'CH001', 'CERT-1001', 'NEW', DATE'2024-01-01', 1200.00, TIMESTAMP'2024-02-05 10:00:00'), 4: ('B0999CH001', 'S1', 'CH001', 'CERT-1001', 'NEW', DATE'2024-01-01', 1250.00, TIMESTAMP'2024-02-12 09:30:00'), -- resubmission 5: ('B0999CH001', 'S1', 'CH001', 'CERT-1002', 'NEW', DATE'2024-01-01', 800.00, TIMESTAMP'2024-02-05 10:00:00'), 6: ('B0999CH001', 'S1', 'CH001', 'CERT-1003', 'NEW', DATE'2024-02-01', 2000.00, TIMESTAMP'2024-03-04 11:00:00'), 7: ('B0999CH001', 'S1', 'CH001', 'CERT-1001', 'MTA', DATE'2024-02-01', 150.00, TIMESTAMP'2024-03-04 11:00:00'), 8: ('B0999CH002', 'S2', 'CH002', 'CERT-2001', 'NEW', DATE'2024-01-01', 5000.00, TIMESTAMP'2024-02-20 15:00:00'), 9: ('B0999CH002', 'S2', 'CH002', 'CERT-2002', 'NEW', DATE'2024-03-01', 3000.00, TIMESTAMP'2024-04-08 15:00:00') 10: AS t(umr, section_no, coverholder_id, certificate_ref, transaction_type, 11: bordereau_month, gross_premium_gbp, submitted_at);
1. Deduplication: keep the latest submission
The classic bordereaux problem. A coverholder corrects a figure and resends the month, and we want only the latest version of each transaction.
1: SELECT umr, section_no, coverholder_id, certificate_ref, transaction_type, 2: bordereau_month, gross_premium_gbp, submitted_at 3: FROM ( 4: SELECT *, 5: ROW_NUMBER() OVER ( 6: PARTITION BY umr, section_no, certificate_ref, transaction_type, bordereau_month 7: ORDER BY submitted_at DESC) AS rn 8: FROM premium_bdx 9: ) x 10: WHERE rn = 1;
Snowflake and Databricks let you write this more neatly with QUALIFY rn = 1. Open-source Apache Spark SQL doesn't support QUALIFY, so on Fabric Spark, EMR or Glue, stick with the subquery. It's portable everywhere.
2. Running totals: GPI utilisation
Every binding authority section has a Gross Premium Income (GPI) limit. Oversight teams want to see cumulative written premium against that limit, month by month.
1: SELECT umr, section_no, bordereau_month, 2: SUM(gross_premium_gbp) AS month_gwp, 3: SUM(SUM(gross_premium_gbp)) OVER ( 4: PARTITION BY umr, section_no 5: ORDER BY bordereau_month 6: ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS cumulative_gwp 7: FROM premium_bdx 8: GROUP BY umr, section_no, bordereau_month;
Join cumulative_gwp to the contract section's GPI limit and you have a utilisation percentage that can drive an alert at 80% and 100%.
Watch the default frame. As in T-SQL, if you put ORDER BY in the window but leave out the frame, the default is RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW. With RANGE, rows that tie on the ORDER BY value are treated as peers and all get the same total. Run it at transaction level, and every transaction in the same bordereau month shows the full month's total. Specify ROWS explicitly unless you really want peer behaviour.
3. Moving windows over a time range
Spark supports RANGE with intervals on date and timestamp columns, which is handy for "claims paid in the last 90 days" calculations where payments are irregular:
1: SELECT claim_ref, payment_date, paid_gbp, 2: SUM(paid_gbp) OVER ( 3: PARTITION BY claim_ref 4: ORDER BY payment_date 5: RANGE BETWEEN INTERVAL 89 DAYS PRECEDING AND CURRENT ROW 6: ) AS paid_last_90_days 7: FROM claim_payments;
T-SQL's RANGE only accepts UNBOUNDED and CURRENT ROW, with no value offsets, so in SQL Server you'd need a self-join or a calendar table for this. One area where Spark is ahead.
4. LAG/LEAD: late bordereaux
Binding authority contracts usually require bordereaux within a set number of days after month end. LAG makes it easy to see each coverholder's submission pattern:
1: WITH submissions AS ( 2: SELECT coverholder_id, bordereau_month, MIN(submitted_at) AS first_submitted 3: FROM premium_bdx 4: GROUP BY coverholder_id, bordereau_month 5: ) 6: SELECT coverholder_id, bordereau_month, first_submitted, 7: DATEDIFF(CAST(first_submitted AS DATE), LAST_DAY(bordereau_month)) AS days_after_month_end, 8: LAG(bordereau_month) OVER (PARTITION BY coverholder_id ORDER BY bordereau_month) AS prev_month_submitted 9: FROM submissions;
Note the argument order: Spark's DATEDIFF(end, start) returns days and takes the end date first, the opposite of T-SQL's DATEDIFF(day, start, end). I've lost more time to that than I want to admit. Newer Spark versions also accept a three-argument DATEDIFF(unit, start, end), but check your runtime before you rely on it.
5. Gaps and islands: missing bordereau months
Which coverholders have a run of consecutive monthly bordereaux, and where are the gaps? The trick is the same as in T-SQL: subtract a row number from the month, and consecutive months end up with the same "island" value.
1: WITH months AS ( 2: SELECT DISTINCT coverholder_id, bordereau_month FROM premium_bdx 3: ), 4: grp AS ( 5: SELECT coverholder_id, bordereau_month, 6: ADD_MONTHS(bordereau_month, 7: -CAST(ROW_NUMBER() OVER (PARTITION BY coverholder_id 8: ORDER BY bordereau_month) AS INT)) AS island 9: FROM months 10: ) 11: SELECT coverholder_id, 12: MIN(bordereau_month) AS run_start, 13: MAX(bordereau_month) AS run_end, 14: COUNT(*) AS months_in_run 15: FROM grp 16: GROUP BY coverholder_id, island 17: ORDER BY coverholder_id, run_start;
In the sample data, CH002 has two separate runs (January, then March), so February's bordereau is missing. That's exactly the kind of thing a DA oversight team wants flagged.
6. Top N per group: largest risks per section
1: SELECT * FROM ( 2: SELECT umr, section_no, certificate_ref, gross_premium_gbp, 3: DENSE_RANK() OVER (PARTITION BY umr, section_no 4: ORDER BY gross_premium_gbp DESC) AS rnk 5: FROM premium_bdx 6: ) WHERE rnk <= 2;
Performance: what's different on a cluster
On SQL Server, a window function might use a sort or an index that's already ordered. On Spark, every PARTITION BY causes a shuffle: rows with the same key are moved to the same executor and then sorted. Some practical points:
- A window without
PARTITION BYputs all rows on a single executor. Spark warns you about this ("No Partition Defined for Window operation"). On years of risk transactions it will be slow or run out of memory. - Skewed keys hurt. If one large coverholder writes 40% of the premium, one task does 40% of the work. Adaptive Query Execution helps with skewed joins, but window skew is still on you. Consider partitioning by contract section rather than coverholder, or splitting out the heavy key.
- Several windows with the same
PARTITION BYandORDER BYshare one shuffle and sort. Windows with different specs each add one. Group your calculations by spec where you can. - Filter early. Restrict to the bordereau months or years of account you need in a CTE before the window, so less data gets shuffled.
The same thing in the DataFrame API
When I'm in Python I often switch to the DataFrame API. The equivalent of the deduplication query:
1: from pyspark.sql import functions as F, Window 2: 3: keys = ["umr", "section_no", "certificate_ref", "transaction_type", "bordereau_month"] 4: w = Window.partitionBy(*keys).orderBy(F.col("submitted_at").desc()) 5: 6: latest = (spark.table("premium_bdx") 7: .withColumn("rn", F.row_number().over(w)) 8: .filter("rn = 1") 9: .drop("rn"))
Both produce the same plan, so use whichever reads better to your team. For SQL people moving to Spark, starting in SQL and moving to DataFrames later is a perfectly good path.
No comments:
Post a Comment