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:
| Worker | DPU | vCPU / memory | Typical use |
|---|---|---|---|
| G.1X | 1 | 4 vCPU / 16 GB | Most ETL; the sensible default |
| G.2X | 2 | 8 vCPU / 32 GB | Memory-heavy joins, wide aggregations |
| G.4X / G.8X | 4 / 8 | 16–32 vCPU / 64–128 GB | Very 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-enable1: 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.
No comments:
Post a Comment