Apache Sqoop moved bulk data between relational databases and Hadoop: tables into HDFS or Hive, and files from HDFS back into tables. For about a decade it was the default answer to "get the nightly extract from Oracle into the cluster". The project was retired in June 2021 and now sits in the Apache Attic, so it receives no fixes, including security fixes. Yet many clusters still run Sqoop jobs that nobody has touched in years, and the people who inherit them need to know exactly what those jobs do before they can safely replace them.

This article explains Sqoop from the mechanism up: what happens on the client, how one table becomes several parallel range queries, where incremental imports silently lose rows, why a failed export leaves half its data in the target table, and how each of those behaviours maps onto a modern replacement. Flags and defaults are from the Sqoop 1.4.7 user guide, the last release line. For where Sqoop sits among other retired Hadoop components, see the Hadoop ecosystem in 2026.

Advertisement

What a Sqoop import actually runs

Sqoop is not a server. It is a command-line client that turns your arguments into a MapReduce job. An import runs in six steps, shown in the diagram below.

First, the client connects to the database over JDBC and reads the table's column names and types. Second, it asks the database for the range of the split column. Third, it generates a Java class that represents one row, with code to read it from a JDBC result set and write it to the chosen file format, and compiles it into a jar. Fourth, it submits a map-only job to YARN with that jar. Fifth, each map task opens its own JDBC connection and runs a query restricted to its slice of the range. Sixth, each mapper writes its rows to one file in the target directory.

Two consequences follow directly. There are no reducers, so the output has exactly one file per mapper, and the file sizes mirror how evenly the rows were split. And each mapper reads in its own database session, so the mappers see the table at slightly different moments: the import is not a consistent snapshot unless the database or connector arranges one. If you need the map side in more detail, the MapReduce map phase covers splits, task lifecycle and map-only jobs.

One sqoop import: metadata on the client, data in parallel map tasksSqoop clientedge nodeRDBMSorders table1. JDBC: column types2. SELECT min(id), max(id)Codegenorders.java + jar3YARNmap-only job4. submitmap 0id 1..10Mmap 1id 10M..20Mmap 2id 20M..30Mmap 3id 30M..40M5. each mapper: own JDBC connection,range query, fetch-size rows per round tripHDFS target dirpart-m-00000 .. part-m-00003 (text, Avro, Parquet or SequenceFile)6. writeNo reducers: one output file per mapper, and no single transaction spans the mappers
Figure 1. A four-mapper import of an orders table. The client does metadata work and code generation; the parallel reads happen in map tasks, each with its own connection and range.

How the split ranges are chosen

Parallelism comes from the split column. By default Sqoop uses the table's primary key and runs select min(<split-by>), max(<split-by>) from <table> to find the boundaries. It then divides that interval evenly into as many ranges as there are mappers, four by default (--num-mappers or -m changes it). Each mapper's query gets a WHERE clause such as id >= 10000001 AND id < 20000001.

The split is even in value, not in row count. If ids 1 to 30 million were bulk-loaded in 2019 and are now mostly deleted, while recent orders sit between 30 and 40 million, three mappers finish almost immediately and one does nearly all the work. You can override the boundaries with --boundary-query, which must return two numeric columns, or choose a different --split-by column that is dense and indexed. An unindexed split column turns every mapper's range query into a full table scan, multiplied by the number of mappers.

If the table has no primary key and you give no --split-by, the import fails unless you set one mapper or pass --autoreset-to-one-mapper, which is mainly meant for import-all-tables across tables of mixed shapes.

Advertisement

A realistic import command

This imports one table into Parquet with eight mappers, reading the password from a protected file rather than the command line, where it would be visible in the process list and shell history:

sqoop import \
  --connect jdbc:postgresql://db01:5432/shop \
  --username etl --password-file /user/etl/.pg.password \
  --table orders \
  --split-by order_id -m 8 \
  --fetch-size 10000 \
  --where "created_at < DATE '2026-10-01'" \
  --as-parquetfile \
  --target-dir /data/raw/shop/orders/dt=2026-09-30

For joins or projections, a free-form query replaces --table. The query must contain the token $CONDITIONS, which each mapper replaces with its own range predicate, and you must give --split-by and --target-dir. Single quotes stop the shell from expanding the dollar sign:

sqoop import \
  --connect jdbc:postgresql://db01:5432/shop \
  --username etl --password-file /user/etl/.pg.password \
  --query 'SELECT o.order_id, o.total, c.country
           FROM orders o JOIN customers c ON c.id = o.customer_id
           WHERE $CONDITIONS' \
  --split-by o.order_id -m 8 \
  --target-dir /data/raw/shop/orders_enriched/dt=2026-09-30

Other flags that matter in practice: --null-string and --null-non-string set how NULL is written in text output (get these wrong and Hive reads the literal string "null"), --direct asks for a database-specific fast path such as a native dump tool where a connector supports it, and --hive-import creates and loads a Hive table after the files land. Other options are -P to prompt at the console, or --password-alias to read from a Hadoop credential provider.

Incremental imports and where they lose rows

Sqoop has two incremental modes. In append mode you name a --check-column whose values only grow, such as an identity id, and Sqoop imports rows where it exceeds --last-value. In lastmodified mode the check column is a timestamp updated on every change, and Sqoop imports rows newer than the last value. At the end of a run Sqoop prints the new last value; if you run the import as a saved job (sqoop job --create, then sqoop job --exec), the value is stored in the job and advanced automatically.

sqoop job --create orders_incr -- import \
  --connect jdbc:postgresql://db01:5432/shop \
  --username etl --password-file /user/etl/.pg.password \
  --table orders --incremental lastmodified \
  --check-column updated_at --last-value '2026-09-30 00:00:00' \
  --merge-key order_id \
  --target-dir /data/raw/shop/orders_current

sqoop job --exec orders_incr      # run nightly; last-value advances itself

With lastmodified, --merge-key folds updated rows into the existing dataset so each key appears once. Both modes have gaps you should know before trusting one:

  • Deletes are invisible. A deleted row has no check-column value to compare, so it stays in the copy forever.
  • Late commits are skipped. A transaction that set updated_at to 01:59 but committed at 02:01, after a 02:00 import recorded 02:00 as its last value, is never imported. Sequences behave the same way in append mode: ids are allocated in order but committed out of order. The usual workaround is to re-read an overlap window, which then needs the merge step to de-duplicate.
  • Updates in append mode are lost. Append only sees new ids; changes to old rows never arrive.
  • The metastore holds the state. If the shared metastore or the local job store is lost, so is every job's last value. The metastore does not store passwords unless sqoop.metastore.client.record.password is set to true, which Oozie-run jobs typically needed, and it is not a secure store.

Exports are not transactions

Export reads files from HDFS and inserts them into a table. Each mapper is a separate writer with its own connection and its own transactions. The user guide is specific: Sqoop uses multi-row INSERT statements of up to 100 records, and commits every 100 statements, so roughly every 10,000 rows per writer. The guide states plainly that an export is therefore not atomic and partial results become visible before it finishes.

If a mapper fails halfway, the rows it already committed stay. Rerunning the job inserts them again, which either fails on a primary key collision or duplicates data in tables without one. Common failure causes are lost connectivity, a row that violates a constraint, and a malformed record in the input files.

The mitigation is --staging-table: mappers write into an empty table with the same structure, and Sqoop moves the staged rows into the target in a single transaction at the end. You create the staging table yourself, and it must be empty or you must pass --clear-staging-table. Staging is not available with --update-key, with stored-procedure exports, and not always with --direct. For updates, --update-key order_id turns inserts into UPDATE statements, and --update-mode allowinsert makes them upserts on connectors that support it. Upserts are idempotent, which is often a better fix than staging: a rerun converges to the same table.

Worked example: sizing a nightly import

Assume an orders table of 40 million rows averaging 200 bytes, about 8 GB, read over a network path that sustains roughly 50 MB/s per JDBC stream, and a database that comfortably serves eight concurrent scans during the batch window.

With the default four mappers and evenly spread ids, each mapper reads 2 GB at 50 MB/s, about 40 seconds of transfer plus query start-up and file writing: a few minutes end to end, dominated by YARN start-up and codegen rather than data. Raise the count to eight and transfer halves, but the database now runs eight range scans at once. At 32 mappers you have usually hit the database's I/O ceiling, and each extra session slows all the others and the production workload alongside them. Size mappers to what the source can serve, not to what the cluster can run.

Now suppose the ids are skewed as described earlier, with 75 percent of rows in the last quarter of the range. One mapper reads 6 GB while three read under 1 GB each, and the job takes about three times longer than the even case regardless of mapper count. A --boundary-query restricted to live rows, or a split column such as a dense surrogate key, restores balance. Finally, check the output: eight Parquet files per day is fine, but per-hour imports of small tables into many partitions create the problem described in the HDFS small files problem.

Failure modes in inherited Sqoop jobs

  • Credentials in plain sight. --password on the command line leaks through process listings, YARN logs and shell history. Search scheduler definitions for it first.
  • Type mapping surprises. Database DECIMAL, DATE and timestamp types map through Java types; time zones and precision can shift between source and file. Compare a sample of values, not just row counts.
  • Inconsistent reads. Mappers read at different moments, so a table under active writes can yield an import that never existed as a single state, such as an order line without its order.
  • Hot-spotting the source. Unindexed split columns or too many mappers degrade the production database during the batch window.
  • Half-written exports. Reruns after a failed insert-only export duplicate data or fail on keys.
  • Unpatched software. Retired means no fixes for Sqoop or its bundled dependencies, which security reviews will eventually flag.

Moving each job off Sqoop

Most Sqoop imports translate directly into Spark's JDBC source, which uses the same idea: a numeric partition column, lower and upper bounds, and a partition count that becomes the number of parallel range queries. Unlike Sqoop, Spark does not query min and max for you, so you run that query yourself (and the bounds only shape the ranges; rows outside them still land in the first and last partitions).

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("orders_extract").getOrCreate()
url = "jdbc:postgresql://db01:5432/shop"
props = {"user": "etl", "password": read_secret("pg/etl"), "fetchsize": "10000"}

lo, hi = (spark.read.jdbc(url, "(SELECT min(order_id) lo, max(order_id) hi FROM orders) b",
                          properties=props).first())

orders = spark.read.jdbc(url, "orders", column="order_id",
                         lowerBound=lo, upperBound=hi + 1, numPartitions=8,
                         properties=props)
(orders.where("created_at < DATE '2026-10-01'")
       .write.mode("overwrite").parquet("/data/raw/shop/orders/dt=2026-09-30"))
SqoopReplacement
--split-by, -mSpark JDBC partitionColumn, numPartitions
--boundary-queryYour own min/max query feeding lowerBound/upperBound
--incremental lastmodifiedChange data capture (for example Debezium into Kafka), which also sees deletes and late commits
Saved job last-valueState in the orchestrator or in the target table's max watermark
--hive-importWrite to a Hive or Iceberg table directly from Spark
Export with stagingSpark JDBC write into a staging table, then one SQL MERGE or INSERT ... SELECT in the database

Incremental loads are the place to improve rather than translate. Log-based change data capture reads the database's transaction log, so it captures deletes, every intermediate update and late commits in commit order, which closes the gaps described above. Writing into Hive or Iceberg tables from Spark then replaces the separate Hive import step. For how the old job ran on the cluster, the MapReduce overview explains the job lifecycle you will see in YARN logs while you migrate.

Trade-offs

ApproachGainCost
Keep SqoopNo migration work nowUnpatched dependency, MapReduce-only, no deletes in incrementals
Spark JDBC batchSame range-split model, modern formats and tablesYou own boundaries, retries and idempotent writes
Log-based CDCDeletes, ordering and low latencyKafka and connector operations, log retention on the source
Database-native exportFast, consistent snapshotVendor-specific tooling and formats

What to do next

  1. Inventory every Sqoop invocation in schedulers, cron and Oozie: table, mode, mapper count, target and owner.
  2. Remove any --password argument; use a password file or credential alias until the job is migrated.
  3. For each incremental job, decide whether deletes or late commits matter. If they do, plan CDC, not a translated batch job.
  4. Check split columns: indexed, dense, and with a sensible boundary query; measure per-mapper output sizes for skew.
  5. Make exports idempotent with upserts, or stage them and commit with one database-side statement.
  6. Port one import to Spark JDBC, run both side by side for a week, and compare row counts and sampled values.
  7. Delete the Sqoop job and its saved-job state only after the replacement has run cleanly through a month end.
Key takeaway: Sqoop is a client that turns a table into a map-only MapReduce job: it reads metadata, generates a row class, splits the min-to-max range of a column evenly across mappers, and has each mapper run its own range query into one output file. That design explains its behaviour: skew follows value gaps, imports are not consistent snapshots, incrementals miss deletes and late commits, and exports commit in chunks and are not atomic unless staged. Sqoop has been retired since 2021, so the useful work now is to understand each job precisely and replace it with Spark JDBC for batch copies and log-based change data capture for incrementals.