June 11, 2026

Tuning AWS Glue PySpark Jobs for Speed and Cost

AWS Glue is easy to start with: write a PySpark script, choose a worker type, click run. It's also easy to overpay for. Most Glue jobs I've reviewed were either over-provisioned "just in case" or spent most of their time on problems that more workers can't fix: small files, reading partitions they didn't need, or skewed joins.

Delegated authority data hits all three. The bordereaux themselves are processed by the bordereaux platform. Downstream, the data team receives the platform's output (its gold layer, or a regular feed of risk, premium and claims transactions) and builds data products on top. Those feeds arrive as many small incremental files. Getting the latest position of each risk tempts jobs to re-read years of history. And a handful of large binders account for most of the premium. This is the checklist I use when a Glue job is slow, expensive, or both.

Understand the bill first

Glue Spark jobs are billed per DPU-hour, per second, with a 1-minute minimum on current Glue versions. Worker types map to DPUs:

WorkerDPUvCPU / memoryTypical use
G.1X14 vCPU / 16 GBMost ETL; the sensible default
G.2X28 vCPU / 32 GBMemory-heavy joins, wide aggregations
G.4X / G.8X4 / 816–32 vCPU / 64–128 GBVery large shuffles; rarely needed

So 10 × G.2X for one hour is 20 DPU-hours, the same as 20 × G.1X. The question is never "how big?" on its own. It's "how big, for how long?"

1. Use the latest Glue version

Glue 5.0 runs Spark 3.5 with Python 3.11, and each new version has brought real performance improvements from newer Spark (Adaptive Query Execution improvements, faster Parquet reading) and faster start-up. Moving older jobs from Glue 3.0 to a current version often makes them noticeably cheaper, with no code changes beyond fixing deprecated APIs.

2. Turn on auto scaling

1: --enable-auto-scaling  true

With auto scaling on, Number of workers becomes a maximum, and Glue adds or removes executors as the job's stages need them. A typical job on the bordereaux feed has one heavy stage (joining the latest transactions to contract sections and converting currencies) and then a long tail of light work, and auto scaling stops you paying for idle executors during that tail. This is the single easiest saving.

3. Use Flex for anything that isn't urgent

The Flex execution class runs jobs on spare capacity at a significantly lower DPU rate (roughly a third cheaper). The trade-off is that start-up can be delayed and runs aren't guaranteed to start immediately.

1: aws glue start-job-run \
2:   --job-name da_claims_bdx_history_rebuild \
3:   --execution-class FLEX

Historical rebuilds, reloading full history after the platform redelivers a feed, and test runs are good candidates. The month-end premium and claims loads that underwriters are waiting for aren't. A simple rule is to split jobs into "must finish for the month-end pack" (Standard) and "anytime overnight" (Flex), and often more than half can go on Flex.

4. Read less data: pushdown predicates

The most common waste I find: a job reads an entire Glue Catalog table, every feed delivery since the platform went live, and then filters to the latest delivery in Spark. Push the filter into the read so Glue only lists and reads the partitions you need:

1: feed_date = args["feed_date"]               # e.g. '2026-05-31', passed as a job argument
2: 
3: premium = glueContext.create_dynamic_frame.from_catalog(
4:     database="da_raw",
5:     table_name="bdx_feed_premium",
6:     push_down_predicate=f"feed_date = '{feed_date}'",
7:     additional_options={"catalogPartitionPredicate": f"feed_date = '{feed_date}'"},
8: )

catalogPartitionPredicate filters on the Glue Catalog side using partition indexes. On feed tables partitioned by delivery date, the partition count grows quickly, and this can cut minutes from the job's start before Spark does any real work.

If you use the Spark DataFrame API directly (spark.read.parquet(...) or spark.table(...)), a WHERE on partition columns gives you partition pruning. Just make sure it's applied before any operation that stops Spark pushing it down.

5. Fix small files on the way in

Incremental feeds from the bordereaux platform tend to arrive as many small files: one per entity, per delivery, sometimes per batch. Over a few years that's a lot of small files to list and open. Glue can group them into larger read tasks:

 1: risk_bdx = glueContext.create_dynamic_frame.from_options(
 2:     connection_type="s3",
 3:     connection_options={
 4:         "paths": ["s3://da-bdx-feed/risk_transaction/"],
 5:         "recurse": True,
 6:         "groupFiles": "inPartition",
 7:         "groupSize": "134217728",       # ~128 MB per group
 8:     },
 9:     format="json",
10: )

6. ...and on the way out

Don't create the next job's small-file problem. Control output file counts:

1: (risk_df.repartition("year_of_account", "bordereau_month")
2:     .write.mode("overwrite")
3:     .partitionBy("year_of_account", "bordereau_month")
4:     .option("maxRecordsPerFile", 5_000_000)
5:     .parquet("s3://da-lake/silver/risk_transaction/"))

Better still, write to an Iceberg table, which Glue supports natively (--datalake-formats iceberg), and let compaction (Glue Data Catalog's automatic optimisation or a scheduled rewrite_data_files) handle file sizes over time. Iceberg's MERGE also makes corrected feed deliveries far easier to apply than overwriting partitions.

7. Job bookmarks for incremental processing

For jobs that process "new feed files since last run", Glue job bookmarks track what has been processed:

1: --job-bookmark-option  job-bookmark-enable
1: src = glueContext.create_dynamic_frame.from_catalog(
2:     database="da_raw", table_name="bdx_feed_claims",
3:     transformation_ctx="src_bdx_feed_claims",  # required for bookmarks to work
4: )
5: # ... transforms ...
6: job.commit()                                 # bookmark only advances on commit

Two gotchas: forget transformation_ctx and the bookmark silently does nothing; forget job.commit() and every run reprocesses everything. And if the platform redelivers a corrected file under the same name, a bookmark may skip it, so agree versioned file names with the feed owner.

8. Look at the metrics before adding workers

Turn on the Spark UI and continuous logging (--enable-spark-ui true, --spark-event-logs-path s3://...) and Glue job observability metrics. I tune Glue jobs from their CloudWatch metrics first, then the Spark UI. Before scaling up, look for:

  • One task taking much longer than the rest in a stage → data skew, often one large binder or facility. More workers won't help. Check AQE skew join handling (spark.sql.adaptive.skewJoin.enabled, on by default), broadcast the smaller side, or salt the hot key.
  • Lots of disk spill → memory pressure. Moving from G.1X to G.2X with fewer workers can be cheaper than adding more G.1X.
  • Executors idle most of the time → over-provisioned. Lower the max workers or rely on auto scaling.
  • The driver busy while executors wait → something is calling collect(), toPandas() or a Python loop over binders on the driver.

9. Broadcast small dimensions

Contract sections, coverholders, brokers and month-end rates of exchange are small next to the transaction facts. Broadcast them:

1: from pyspark.sql.functions import broadcast
2: 
3: premium_enriched = (premium_df
4:     .join(broadcast(dim_contract_section), ["umr", "section_no"], "left")
5:     .join(broadcast(dim_broker), "broker_code", "left")
6:     .join(broadcast(month_roe), ["original_ccy", "bordereau_month"], "left"))

AQE often does this automatically when it can see the size, but an explicit hint is reliable when the dimension comes from a source Spark can't size well, such as a JDBC read from an in-house SQL Server.

Results

Applying this list (Glue version upgrade, auto scaling, Flex for non-urgent jobs, pushdown predicates, and output compaction) usually saves a lot, and the month-end jobs often finish earlier too. None of it is clever. It's just checking the basics one job at a time.

Do the cheap things first, measure with the Spark UI, and only then pay for more workers.

January 22, 2026

Spark 4.0 for SQL Developers: ANSI Mode, VARIANT, Pipe Syntax and Collations

Apache Spark 4.0 came out last year, and it's now arriving in managed platforms' runtimes. If you're about to move your workloads from a Spark 3.x runtime, there are a few SQL changes you need to know about. One of them can break existing jobs, and bordereaux pipelines are particularly exposed to it. The rest are good additions, especially for people who think in SQL first.

Here are the four SQL changes I think matter most, with delegated authority examples you can try in a Spark 4 session.

1. ANSI mode is now on by default

This is the one that can break things, so it comes first.

In Spark 3.x, spark.sql.ansi.enabled was false by default, and Spark was forgiving in ways that hid bad data. Having mapped dozens of different bordereaux formats over the years, I've seen every variation of the values that turn up in a premium column:

1: -- Spark 3.x (non-ANSI)
2: SELECT CAST('1,250.00' AS DECIMAL(18,2));   -- NULL (thousands separator)
3: SELECT CAST('N/A' AS DECIMAL(18,2));        -- NULL
4: SELECT CAST('31/02/2025' AS DATE);          -- NULL
5: SELECT 250.00 / 0;                          -- NULL (commission % with zero premium)

In Spark 4.0 with ANSI mode on, all of these raise errors. As a former SQL Server developer I think that's correct. SQL Server has always failed on these, and quietly turning a premium of "1,250.00" into NULL meant written premium was understated with nobody noticing. But a pipeline that "worked" on 3.x can start failing on 4.0 the first time a coverholder sends a bad value.

How to handle it:

 1: -- Clean known formats, then use TRY_ functions where bad values are expected
 2: SELECT certificate_ref,
 3:        TRY_CAST(REPLACE(gross_premium_text, ',', '') AS DECIMAL(18,2)) AS gross_premium,
 4:        CAST(TRY_TO_TIMESTAMP(inception_date_text, 'dd/MM/yyyy') AS DATE) AS inception_date
 5: FROM   bronze.premium_bdx_raw;
 6: 
 7: SELECT umr, section_no,
 8:        TRY_DIVIDE(commission_amount, gross_premium) * 100 AS commission_pct
 9: FROM   silver.premium_transaction;
10: 
11: -- Find the offending rows, by coverholder, before migrating
12: SELECT coverholder_id, gross_premium_text, COUNT(*) AS rows_affected
13: FROM   bronze.premium_bdx_raw
14: WHERE  gross_premium_text IS NOT NULL
15:   AND  TRY_CAST(REPLACE(gross_premium_text, ',', '') AS DECIMAL(18,2)) IS NULL
16: GROUP  BY coverholder_id, gross_premium_text
17: ORDER  BY rows_affected DESC;

That last query is worth turning into a standard data quality report that goes back to coverholders. You can set spark.sql.ansi.enabled = false to get the old behaviour back, and that's a reasonable short-term fix to unblock a migration. I'd treat it as temporary, though. The errors are pointing at real data problems.

2. The VARIANT type for semi-structured data

Before 4.0, JSON in Spark meant either keeping it as a STRING and parsing it with get_json_object/from_json on every query, or defining a fixed STRUCT schema that broke when the source added a field. Snowflake users have had VARIANT for years. Now Spark has it too.

It's a good fit for bordereaux submitted through APIs or portals, where each coverholder's payload carries the standard fields plus their own extras:

 1: CREATE TABLE bronze.bdx_submission (
 2:   submission_id   BIGINT,
 3:   coverholder_id  STRING,
 4:   received_at     TIMESTAMP,
 5:   payload         VARIANT
 6: ) USING DELTA;
 7: 
 8: INSERT INTO bronze.bdx_submission
 9: SELECT submission_id, coverholder_id, received_at, PARSE_JSON(raw_json)
10: FROM   landing.bdx_submission_raw;
11: 
12: -- Extract with a path and a target type
13: SELECT submission_id,
14:        VARIANT_GET(payload, '$.contract.umr', 'STRING')               AS umr,
15:        VARIANT_GET(payload, '$.contract.section', 'STRING')           AS section_no,
16:        VARIANT_GET(payload, '$.risk.certificate_ref', 'STRING')       AS certificate_ref,
17:        VARIANT_GET(payload, '$.risk.sum_insured', 'DECIMAL(18,2)')    AS sum_insured,
18:        VARIANT_GET(payload, '$.premium.gross', 'DECIMAL(18,2)')       AS gross_premium,
19:        TRY_VARIANT_GET(payload, '$.risk.flood_zone', 'STRING')        AS flood_zone  -- only some coverholders send it
20: FROM   bronze.bdx_submission;
21: 
22: -- See what a payload looks like
23: SELECT SCHEMA_OF_VARIANT(payload) FROM bronze.bdx_submission LIMIT 1;

VARIANT is stored in a binary format, so it's much faster to query than re-parsing JSON strings, and it handles schema changes without DDL. Note that table formats have to support it too. Delta Lake added VARIANT support alongside Spark 4.0, so check the Delta version on your platform before relying on it in shared tables.

3. SQL pipe syntax

This one divides opinion, but I've grown to like it. Pipe syntax lets you write a query in the order it runs, using the |> operator:

1: FROM silver.premium_transaction
2: |> WHERE year_of_account = 2025
3: |> AGGREGATE SUM(gross_premium_gbp) AS gwp_gbp, COUNT(*) AS transactions
4:    GROUP BY coverholder_id
5: |> WHERE gwp_gbp > 100000
6: |> ORDER BY gwp_gbp DESC
7: |> LIMIT 10;

Compare that with the usual SELECT ... FROM ... WHERE ... GROUP BY ... HAVING ... ORDER BY, where the order you write the clauses isn't the order they're evaluated. With pipes:

  • There's no HAVING to remember. It's just another WHERE after the aggregation.
  • You can add steps without wrapping everything in another subquery or CTE.
  • It reads like a DataFrame chain, which helps when the same team writes both SQL and PySpark.

A longer example: GPI utilisation by contract section, with a join and derived columns:

1: FROM silver.premium_transaction AS p
2: |> JOIN silver.contract_section AS cs USING (umr, section_no)
3: |> WHERE cs.year_of_account = 2025
4: |> AGGREGATE SUM(p.gross_premium_gbp) AS gwp_gbp, MAX(cs.gpi_limit_gbp) AS gpi_limit_gbp
5:    GROUP BY umr, section_no, cs.coverholder_id
6: |> EXTEND ROUND(gwp_gbp / gpi_limit_gbp * 100, 1) AS gpi_utilisation_pct
7: |> WHERE gpi_utilisation_pct >= 80
8: |> ORDER BY gpi_utilisation_pct DESC;

It's fully optional and mixes freely with regular SQL. I wouldn't rewrite existing code, but for new exploratory queries it's become my habit.

4. String collations

Another feature SQL Server developers have always had and Spark lacked: collations. Coverholder and broker names arrive in every combination of upper and lower case, depending on who keyed them. You can now compare and group strings case-insensitively without wrapping everything in LOWER():

 1: SELECT 'Acme Underwriting Ltd' = 'ACME UNDERWRITING LTD' COLLATE UTF8_LCASE;   -- true
 2: 
 3: CREATE TABLE silver.broker (
 4:   broker_code   STRING,
 5:   broker_name   STRING COLLATE UTF8_LCASE,
 6:   broker_group  STRING COLLATE UNICODE_CI
 7: ) USING DELTA;
 8: 
 9: SELECT broker_name, COUNT(*)
10: FROM   silver.broker
11: GROUP  BY broker_name;     -- 'Acme Re Brokers' and 'ACME RE BROKERS' group together

Like VARIANT, collations in stored tables depend on table format support, so test on your platform first. Expressions using COLLATE in a query work regardless. And collations only fix case and accents. "Acme Underwriting Ltd" vs "Acme Underwriting Limited" still needs proper reference data matching on the coverholder PIN or broker code.

Other things worth a look

  • SQL scripting: BEGIN ... END blocks with variables, IF, WHILE and loops. It's familiar territory for anyone who wrote T-SQL stored procedures. It started as a preview in 4.0, so check its status on your runtime.
  • Python data source API: write custom readers and writers in pure Python, which is useful for reading from bordereaux platform APIs or niche file formats.
  • Spark Connect improvements: a thin client architecture that more platforms now build on.

A migration checklist

  1. Run your existing jobs on a Spark 4 runtime in a test environment with ANSI on, against several months of real bordereaux. Collect every failure.
  2. Fix casts and divisions with TRY_ functions where bad data is expected, and push data quality issues back to coverholders where it isn't.
  3. Check third-party libraries and connectors for Spark 4 / Scala 2.13 builds. Spark 4.0 dropped Scala 2.12 and Java 8/11, and requires Java 17 or later.
  4. Only then start using VARIANT, collations and pipe syntax in new code.

Spark 4.0 brings Spark SQL closer to what SQL developers expect from a database: strict about bad data, flexible with semi-structured data, and more pleasant to write. For anyone processing bordereaux, the strictness alone is worth the upgrade.