Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Reading from JDBC sources in parallel

PySpark · DataFrame API in Practice

Reading from JDBC sources in parallel

Mediumpyspark-80
jdbcparallel-readpartitioncolumndatabase

Question

How do you read a large table from Postgres or MySQL into Spark without a single slow task?

Solution

By default, a JDBC read in Spark uses one connection and one partition, so one task pulls the whole table and everything else waits. To parallelise, you give Spark a numeric (or date) column and a range, and it opens several connections that each read a slice.

The four options

df = (spark.read.format("jdbc")
      .option("url", "jdbc:postgresql://replica-host:5432/shop")
      .option("dbtable", "public.orders")
      .option("user", user).option("password", password)
      .option("partitionColumn", "order_id")
      .option("lowerBound", 1)
      .option("upperBound", 100_000_000)
      .option("numPartitions", 16)
      .option("fetchsize", 10_000)
      .load())

Spark splits the range 1 to 100,000,000 into 16 strides and issues 16 queries such as WHERE order_id >= 6250001 AND order_id < 12500001.

The bounds do not filter

lowerBound and upperBound only decide how to cut the strides. Rows below the lower bound go into the first partition, and rows above the upper bound go into the last. They are still read. To restrict the data, add a real filter, for example through a subquery in dbtable.

Even slices

Pick a column with evenly spread values, such as an auto-increment id. If the column is skewed (many rows near one end), some partitions are huge. A date column works if data is spread across days. A string column does not work with this method.

Pushdown and the query option

Filters you apply on the DataFrame are pushed down into the SQL where possible, which you can see in explain(). To send a custom query, use dbtable with a subquery like (SELECT ... ) AS t, or the query option. The query option cannot be combined with partitionColumn, so use the subquery form when you need partitioned reads.

Be kind to the source

Sixteen parallel connections on a production OLTP database can slow down the application. Read from a read replica, keep numPartitions modest (often 4 to 16), and run large extracts off-peak. fetchsize controls how many rows come per round trip. The default for some drivers is very low, and raising it is a cheap speedup. If the pipeline needs to run often, consider CDC instead of repeated full table reads.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext