Plenty of Redshift work never touches a COPY statement. Teams that already run Apache Spark for machine learning feature engineering, Iceberg maintenance or heavy file-level transformation want Redshift as both a source and a sink, and they want it without hand-writing UNLOAD/COPY plumbing around every job.
That is what the Amazon Redshift integration for Apache Spark is for. It ships built in with AWS Glue 4.0+ and Amazon EMR 6.9+, so there is no connector JAR to chase, no spark-redshift fork to maintain, and no JDBC driver version roulette. This guide covers the wiring, what actually gets pushed down to Redshift, how writes behave, and the traps that turn a five-minute job into a forty-minute one.
How the integration actually moves data
The mental model matters more than the syntax. The connector uses two channels:
- A JDBC control channel. Spark opens a JDBC session to Redshift to run metadata queries, issue
UNLOADfor reads andCOPYfor writes, and to run the transactional DDL around a write. - A bulk channel through S3. Data does not stream through the JDBC connection. On read, Redshift
UNLOADs the result set to a staging prefix in S3 in Parquet, and Spark executors read those files in parallel. On write, Spark writes staging files to S3 and RedshiftCOPYs them.
Everything that goes wrong with this connector goes wrong in one of those two places: the JDBC session (permissions, network, timeouts) or the S3 staging prefix (tempdir, encryption, lifecycle, region).
Prerequisites and IAM
You need three identities lined up:
- The Spark job role (the Glue job role or the EMR EC2 instance profile) needs
s3:GetObject,s3:PutObject,s3:DeleteObjectands3:ListBucketon the staging prefix, plus permission to read the Secrets Manager secret holding the Redshift credentials. - A Redshift IAM role associated with the cluster or Serverless namespace, used by
UNLOAD/COPYto reach the same S3 prefix. This is passed asaws_iam_role. - A Redshift database user with
SELECTon the source tables and, for writes,CREATE/INSERTon the target schema plus a schema it can create temporary staging tables in.
Keep the staging bucket in the same region as the Redshift namespace. Cross-region UNLOAD either fails outright or silently adds transfer cost and latency.
Network-wise, Glue and EMR must reach Redshift's VPC endpoint (security group ingress on 5439 from the job's subnet/security group) and must reach S3 — a gateway VPC endpoint for S3 in the job's route table, or the job will stall on a private subnet with no NAT.
Reading from Redshift
Python, Glue or EMR alike:
secret = "arn:aws:secretsmanager:us-east-1:111111111111:secret:redshift/analytics"
opts = {
"url": "jdbc:redshift://prod.123456789012.us-east-1.redshift-serverless.amazonaws.com:5439/analytics",
"tempdir": "s3://acme-spark-staging/redshift/",
"aws_iam_role": "arn:aws:iam::111111111111:role/RedshiftSpectrumAndUnloadRole",
"secret.id": secret,
"tempformat": "PARQUET",
}
orders = (
spark.read.format("io.github.spark_redshift_community.spark.redshift")
.options(**opts)
.option("dbtable", "sales.fact_orders")
.load()
)
Use .option("query", "...") instead of dbtable when you want Redshift to do the heavy lifting explicitly:
q = """
SELECT customer_id, DATE_TRUNC('month', order_ts) AS month, SUM(net_amount) AS revenue
FROM sales.fact_orders
WHERE order_ts >= DATEADD(month, -13, CURRENT_DATE)
GROUP BY 1, 2
"""
monthly = spark.read.format("io.github.spark_redshift_community.spark.redshift").options(**opts).option("query", q).load()
On Glue you can do the same through the catalog with glueContext.create_dynamic_frame.from_options(connection_type="redshift", ...), which resolves the connection and secret for you. Underneath, it is the same UNLOAD-to-S3 mechanism.
Pushdown: what Redshift does versus what Spark does
The integration applies query pushdown: it rewrites parts of the Spark logical plan into the SQL sent to Redshift, so the UNLOAD only materialises what Spark actually needs. In current Glue/EMR releases the pushdown covers:
- projections (column pruning)
- filters and most scalar expressions
- sorts and limits
- aggregations (
SUM,COUNT,AVG,MIN,MAX,GROUP BY,DISTINCT) - joins between two Redshift tables read through the same connection
What does not push down: joins between a Redshift DataFrame and a non-Redshift DataFrame (an S3 Parquet table, a Kafka batch, a Glue catalog table backed by Iceberg), window functions in some shapes, and anything wrapped in a Python UDF. Those force a full unload of whatever the connector could not eliminate.
Verify rather than assume. Print the plan:
monthly.explain(True)
and look for the RedshiftScan/RedshiftRelation node — the pushed SQL text appears in it. Cross-check on the Redshift side:
SELECT query_id, query_text, start_time, elapsed_time
FROM SYS_QUERY_HISTORY
WHERE query_text ILIKE '%UNLOAD%'
AND start_time > DATEADD(hour, -2, GETDATE())
ORDER BY start_time DESC;
If the unloaded SQL is SELECT * FROM sales.fact_orders when your job only wanted thirteen months of three columns, you have lost the pushdown. The usual causes are a UDF applied too early, a cache() between the read and the filter, or a join against a non-Redshift source that you can fix by filtering the Redshift side explicitly first.
Writing back to Redshift
(
features.write.format("io.github.spark_redshift_community.spark.redshift")
.options(**opts)
.option("dbtable", "ml.customer_features")
.option("tempformat", "PARQUET")
.mode("append")
.save()
)
The modes behave differently enough to matter:
| Mode | Behaviour |
|---|---|
append | COPY into the existing table. Table must exist or is created from the Spark schema. |
overwrite | Drops/truncates and reloads the target inside a transaction. Default behaviour drops the table, which loses your DISTKEY, SORTKEY, encodings and grants. |
errorifexists / ignore | Rarely what you want in a pipeline. |
Two options save you from the overwrite footgun:
.option("usestagingtable", "true")— the default; data lands in a staging table and is swapped in, so readers never see an empty target..option("preactions", "...")and.option("postactions", "...")— arbitrary SQL run in the same transaction, before and after theCOPY.
For a managed table you care about, prefer creating the DDL yourself, loading to a staging table, and doing the merge in Redshift SQL:
(
features.write.format("io.github.spark_redshift_community.spark.redshift")
.options(**opts)
.option("dbtable", "ml.customer_features_stg")
.option("preactions", "TRUNCATE TABLE ml.customer_features_stg;")
.option("postactions", """
MERGE INTO ml.customer_features t
USING ml.customer_features_stg s ON t.customer_id = s.customer_id
WHEN MATCHED THEN UPDATE SET feature_vec = s.feature_vec, updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT VALUES (s.customer_id, s.feature_vec, s.updated_at);
""")
.mode("append")
.save()
)
That keeps table design under your control and keeps the write atomic from a reader's point of view.
Performance notes that actually move the needle
Set tempformat to PARQUET. The legacy default of CSV/AVRO is meaningfully slower on both sides and mangles some types. Parquet also preserves numeric precision instead of round-tripping through text.
Control parallelism at the source. A single huge unload file limits read parallelism. UNLOAD writes one file per slice by default, which is usually right; if you force MAXFILESIZE too high you will end up with fewer, larger files and idle executors.
Do not use Spark as a row-by-row writer. If your job produces a few thousand rows, the S3+COPY round trip is pure overhead — write with the Redshift Data API instead.
Push the filter before the join. The single most common performance bug we see is a Redshift DataFrame joined to a lake table with no predicate on the Redshift side, unloading a full fact table every run.
Watch concurrency. Ten Glue jobs each unloading a fact table will queue against your WLM configuration just like any other workload. Give ETL-style unloads their own query priority, or isolate them on a separate consumer warehouse via data sharing.
Temp-directory hygiene
The staging prefix is not self-cleaning in every configuration, and it accumulates copies of production data.
- Put an S3 lifecycle rule on the staging prefix expiring objects after 1–3 days.
- Enable SSE-KMS on the bucket and grant the Redshift IAM role
kms:Decrypt/kms:GenerateDataKey; otherwiseUNLOADfails with an opaque access-denied. - Give the staging prefix its own bucket or prefix per environment. Sharing a
tempdirbetween prod and dev jobs is a data-classification incident waiting for an auditor. - Treat the prefix as containing the same data classification as the source tables — because it does.
Failure patterns and what they mean
| Symptom | Usual cause |
|---|---|
S3ServiceException: Access Denied during read | Redshift IAM role lacks access to tempdir, or KMS key policy omits it |
| Job hangs then times out on a private subnet | Missing S3 gateway VPC endpoint or NAT route |
The specified bucket does not exist region error | Staging bucket in a different region from the Redshift namespace |
| Unexpectedly slow, huge UNLOAD in SYS_QUERY_HISTORY | Lost pushdown — UDF, cache or cross-source join before the filter |
| Target table lost its sort/dist keys | overwrite mode dropping and recreating the table |
| Duplicate rows after a retry | append mode with no idempotency key; use staging table + MERGE |
When not to use Spark here at all
Be honest about the alternative. If the work is SQL-shaped — joins, aggregations, SCD merges over data already in Redshift or reachable via Spectrum, Iceberg or zero-ETL — then dbt or stored procedures running inside Redshift will be faster, cheaper and far easier to operate than a Spark job that pulls data out and pushes it back. Spark earns its place when you need Python/Scala libraries Redshift cannot run (ML feature transforms, custom parsers, image or text processing), when you are already orchestrating a large lake pipeline, or when Redshift is one sink among several.
Where this usually lands in a real engagement
The pattern we see most often is a feature pipeline: Spark on Glue reads dimension and fact data from Redshift, joins it to raw event data in S3, computes features, and writes the result back to a Redshift table that BI and Redshift ML both read. Done carelessly it unloads terabytes a night. Done well — filters pushed down, Parquet tempformat, staging table plus MERGE, lifecycle-expired staging prefix — it moves only what changed.
If you want a second pair of eyes on a Spark-to-Redshift pipeline that has grown expensive, or help deciding what belongs in Spark versus in the warehouse, get in touch. Our senior Redshift consultants do this work on direct engagements and as subcontractors to agencies and consultancies.