Sooner or later every Spark pipeline has to talk to a relational database: pull a customer table out of Postgres, snapshot an Oracle ledger into the lake, or push aggregated results back into MySQL for an application to serve. Spark ships a JDBC data source for this, and it looks deceptively simple: give it a URL and a table name and you get a DataFrame. The defaults then quietly read the whole table through one connection on one executor, or open two hundred connections against a database sized for twenty.
This article explains what the JDBC source actually does: how it splits a read into partitions and the exact WHERE clauses it generates, what it pushes down to the database, how fetch size and type mapping affect memory and correctness, how writes are batched and committed, and why task retries can duplicate rows. It ends with failure modes, a safe upsert pattern and a checklist. Option names and defaults were checked against the Spark 4.2.0 JDBC documentation; most of them have been stable across the 3.x line, but confirm against the docs for your version.
What a JDBC read actually does
A JDBC read has two phases. On the driver, Spark opens a connection, runs a query that returns no rows, and reads the result-set metadata to build the schema. It then computes one WHERE clause per partition. On the executors, each task opens its own connection, runs SELECT <columns> FROM <table> WHERE <partition clause> AND <pushed filters>, and streams rows into Spark's internal format. There is no shared connection pool and no coordination between tasks: each partition is an independent query, in an independent session, that may see a different snapshot of the table.
With only url and dbtable there is one partition, so one task reads every row through one connection. That is fine for a lookup table and terrible for a hundred-million-row fact table: one executor core does all the work while the rest of the cluster idles, and the long-running query holds database resources for the whole duration.
Partitioned reads and the stride
Parallel reads need four options set together: partitionColumn (a numeric, date or timestamp column), lowerBound, upperBound and numPartitions. The documentation is explicit that the bounds only decide the partition stride; they do not filter rows. Spark computes a stride of roughly (upper - lower) / numPartitions, gives the first partition everything below the first boundary plus NULLs, and gives the last partition everything at or above the last boundary.
orders = (spark.read.format("jdbc")
.option("url", "jdbc:postgresql://db.internal:5432/shop")
.option("dbtable", "public.orders")
.option("user", user).option("password", password)
.option("partitionColumn", "id")
.option("lowerBound", "1")
.option("upperBound", "1000001")
.option("numPartitions", "4")
.option("fetchsize", "10000")
.load())
# Generated per task (simplified):
# WHERE id < 250001 OR id IS NULL
# WHERE id >= 250001 AND id < 500001
# WHERE id >= 500001 AND id < 750001
# WHERE id >= 750001Worked example: the orders table has ids from 1 to 1,000,000, but the most recent 600,000 were created this year after a migration that started ids at 400,001, and ids 1 to 400,000 are sparse archive rows. Partition 0 gets a few thousand rows, partition 1 a few tens of thousands, and partitions 2 and 3 get almost everything. The job is as slow as its largest partition. The fix is to choose bounds from the data, not from the schema: run SELECT min(id), max(id) FROM orders first, and if the distribution is uneven, use explicit predicates instead.
If you set the bounds too narrow, nothing is lost, but the edge partitions absorb everything outside them, which is a subtle way to get one giant task. If the column is not indexed, each of the partition queries is a full scan, so four partitions cost four full scans. Partition on the primary key or another indexed column, and check the database's query plan for one of the generated statements.
For full control, pass a list of predicates, one per partition, to the jdbc method. Each string becomes the WHERE clause of one task, so you can partition by date ranges, by hash buckets, or by anything the database can evaluate:
preds = [f"mod(abs(hashtext(customer_id)), 8) = {i}" for i in range(8)]
df = spark.read.jdbc(url, "public.orders", predicates=preds, properties=props)Predicates must be disjoint and together cover every row you want; Spark does not check either property. Overlapping predicates silently duplicate rows and gaps silently drop them, so assert counts against a count(*) on the source for the first runs. The hashtext function is Postgres-specific; use the equivalent in your database.
Pushdown and custom queries
Spark pushes simple filters on the DataFrame into the generated SQL when pushDownPredicate is true, the default. Column pruning is also pushed, so selecting three columns fetches three columns. Anything Spark cannot translate, such as most UDFs, runs in Spark after every row has crossed the network. The documentation also lists pushDownAggregate, pushDownLimit, pushDownOffset and pushDownTableSample, all defaulting to true in 4.2.0 and described as applying to the V2 JDBC path. An aggregate is pushed completely only when numPartitions is 1 or the grouping key is the partition column; otherwise the database computes partial results and Spark finishes them. Do not assume: call explain() and look for the pushed filters and aggregates in the scan node.
When you need a join, window function or anything else Spark would do badly over the wire, let the database run it with the query option. Spark wraps it as a subquery aliased spark_gen_alias. You cannot combine query with dbtable or with partitionColumn; to partition a custom query, put it in dbtable as a parenthesised, aliased subquery. For SQL that cannot live inside a subquery, such as a SQL Server WITH clause, prepareQuery supplies a prefix. The deeper model behind pushdown is in Spark Data Source V2.
Fetch size, types and session setup
fetchsize controls how many rows the driver fetches per network round trip, and its default of 0 means the JDBC driver decides. Driver defaults vary widely: the Oracle driver fetches 10 rows at a time, which makes large reads crawl, while the Postgres driver fetches the whole result into memory unless it is told to use a cursor. Spark's Postgres dialect turns off auto-commit when a positive fetch size is set so that the driver can stream; MySQL Connector/J streams only with useCursorFetch=true in the URL. Start with 1,000 to 10,000 and measure; very wide rows want smaller values.
Type mapping is the second trap. Spark infers types from result-set metadata, which can be lossy: Oracle NUMBER without precision, unsigned MySQL integers, Postgres numeric with large precision and vendor-specific types all deserve a look. customSchema overrides inferred types for named columns, for example id DECIMAL(38, 0), amount DECIMAL(18, 2). For timestamps without a time zone, the default maps them to Spark's session-zone-aware TimestampType; setting preferTimestampNTZ to true maps them to TimestampNTZType and avoids shifts when the Spark session zone differs from the database's intent.
Session setup belongs in sessionInitStatement, which runs once after each connection opens: a statement timeout, a search path, or a read-only transaction. queryTimeout bounds each statement in seconds. Keep credentials out of code and notebooks; Spark secrets management covers the options.
Writes, transactions and duplicates
On write, each DataFrame partition becomes one task with one connection that runs a parameterised INSERT and sends rows in batches of batchsize (default 1,000). numPartitions on a write caps concurrency by coalescing the DataFrame to that many partitions before writing. isolationLevel defaults to READ_UNCOMMITTED. When the driver supports transactions and the level is not NONE, Spark turns off auto-commit and commits once at the end of each partition, so each partition is roughly one transaction and the job as a whole is not.
That has a consequence worth testing in your own environment. If a task fails after its commit, or a speculative duplicate also commits, the rows are inserted twice; if the job fails halfway, some partitions are committed and others are not. Append mode into a table without a unique constraint is therefore at-least-once. Save modes add their own surprises: overwrite drops and recreates the table by default, losing indexes, grants and column types, unless truncate is true, which the documentation says the MySQL, DB2, SQL Server, Derby and Oracle dialects support and Postgres does not. cascadeTruncate truncates dependent tables too, which is rarely what you want.
The robust pattern is a staging table plus a database-side merge. Spark overwrites a staging table it owns, then one statement moves rows into the target idempotently:
(result.repartition(8)
.write.format("jdbc")
.option("url", url).option("dbtable", "staging.daily_totals")
.option("user", user).option("password", password)
.option("truncate", "true") # where the dialect supports it
.option("batchsize", "5000")
.mode("overwrite")
.save())
# Then, in one transaction on the database (Postgres syntax):
# INSERT INTO app.daily_totals AS t (day, store_id, total)
# SELECT day, store_id, total FROM staging.daily_totals
# ON CONFLICT (day, store_id) DO UPDATE SET total = EXCLUDED.total;The merge runs once, inside a real transaction, and is idempotent under replay because of the conflict key. Duplicates from task retries land in staging, where a SELECT DISTINCT or a primary key on staging catches them. Run the merge from the driver with a plain JDBC connection or from your orchestrator after the Spark step succeeds.
Respecting the database
Every partition is a connection. A read with 64 partitions against a Postgres primary that allows 100 connections, while the application holds 60, will fail some tasks with connection refusals and then retry them into the same wall. Size numPartitions from the database's side: how many concurrent scans can it serve without hurting production traffic? Point heavy reads at a replica, schedule them off-peak, and set a statement timeout so a stuck query does not hold locks forever. For repeatable extracts of large tables, consider change data capture or a database-native export instead of re-reading everything nightly.
On the Spark side, the read parallelism you choose for the database is rarely the parallelism you want downstream. Read with eight partitions to protect the database, then repartition for the transformations; see coalesce versus repartition and Spark partitioning. Uneven predicates produce skewed tasks, diagnosed as in Spark data skew.
Failure modes
- Single-partition reads. No partition options, one task, hours of runtime. Check the stage's task count before anything else.
- Bounds mistaken for filters. A job meant to read last week's ids reads the whole table. Put the filter in a DataFrame
whereor in the subquery, and use bounds only for stride. - Skewed partitions. Monotonic ids with gaps or a timestamp column with bursty history create one huge task. Derive bounds from data or use predicates.
- Unindexed partition column. N partitions become N full scans. Partition on an indexed column.
- Executor memory blowups. Fetch size 0 with the Postgres driver buffers a whole partition on the executor and causes out-of-memory errors. Set
fetchsize. - Inconsistent snapshots. Each partition reads at a different moment, so a table under heavy writes yields a mix of states. Read from a snapshot, a replica paused for the extract, or a staging copy when consistency matters.
- Duplicate rows on write. Retries and speculation re-run committed partitions. Write to staging and merge with a key.
- Overwrite drops the table. Indexes, grants and constraints disappear. Use
truncatewhere supported or the staging pattern. - Connection storms. Too many partitions exhaust the database. Cap
numPartitionsand use a replica.
What to do next
- Find every JDBC read in your jobs and record its partition count from the Spark UI.
- For each large table, run
min,maxand a histogram on the candidate partition column, and pick bounds or predicates from the data. - Set
fetchsizeexplicitly and confirm streaming behaviour for your driver. - Run
explain()on each read and confirm filters and columns are pushed down. - Pin risky column types with
customSchemaand decide onpreferTimestampNTZ. - Replace direct appends to serving tables with a staging table and an idempotent merge.
- Agree a connection budget with the database owner and set
numPartitions,sessionInitStatementandqueryTimeoutto respect it. - Add a row-count reconciliation between source and target to every extract.