Skip to content

fix(sql): keep a dedup id on a single job - #31

Open
thetutlage wants to merge 2 commits into
mainfrom
fix/sql-dedup-unique-slot
Open

thetutlage wants to merge 2 commits into
mainfrom
fix/sql-dedup-unique-slot

Conversation

@thetutlage

@thetutlage thetutlage commented Oct 7, 2026 •

Copy link
Copy Markdown

Problem

With the Knex and Kysely adapters, a dispatch using .dedup() could:

  • enqueue a second job for the same dedup id, or
  • return added with the id of a job that was never stored (Knex), or throw Unable to resolve concurrent dedup dispatch (Kysely).

Both happen when the dispatch races with a worker acquiring the job that owns the dedup id.

Cause

The unique index only covered pending and delayed jobs, so a job left the index as soon as it was acquired. A concurrent dispatch that had read no owner then:

  • inserted its job without conflict, since the owner was now active, or
  • lost the insert, then looked for the winner among pending and delayed jobs only, and found none.

MySQL has no partial index, so it had no constraint at all (the caveat documented in the README). Its SELECT ... FOR UPDATE through the dedup index also deadlocked with the completion of the owner.

Fix

Schema. A plain unique index on (queue, dedup_id) replaces the partial one. A single job owns a dedup id, whatever its status. It works the same on PostgreSQL, MySQL and SQLite, because jobs without dedup hold NULL, which is never a duplicate. Removing the job releases the id, so completeJob, failJob and pruning are unchanged.

Dispatch. Still one transaction, with the same two cases:

  • The id has an owner. It is locked with SELECT ... FOR UPDATE, then skipped, extended, replaced, or released when expired, with plain updates under that lock.
  • The id has no owner. There is no row to lock, so the job is inserted and the unique index arbitrates. A dispatch that loses the insert runs its transaction again, and this time finds the winner and locks it.

What changed compared to main:

  • The savepoint, and the lookup of the winner among pending and delayed jobs, are gone. Running the transaction again replaces both.
  • The owner is found through the dedup index without a lock, then locked through its primary key, as a worker locks a job. Locking through the dedup index takes the index entry then the row, while the removal of a completed job takes the row then the index entry. That is the MySQL deadlock.
  • An expired id is released whatever the status of its job. Retained completed and failed jobs used to keep theirs, which the unique index does not allow.
  • On MySQL, InnoDB can pick a dispatch as a deadlock victim when several insert the same id at once. That transaction runs again too.

Upgrading

Existing tables need the new index. addDedupColumns() creates it, and stays idempotent:

// Knex
await new KnexQueueSchemaService(connection).addDedupColumns('queue_jobs')

// Kysely
await new KyselyQueueSchemaService(db, { dialect: 'postgres' }).addDedupColumns('queue_jobs')
  • Where several jobs share a dedup id, the latest keeps it and the others lose only their dedup id. No job is removed.
  • The previous indexes are dropped.
  • Dispatches with .dedup() should be stopped while it runs.

Without the migration, the new code still deduplicates sequential dispatches, but concurrent ones can race as before.

Tests

  • dedup: concurrent dispatches stay deduplicated while jobs are acquired: 4 dispatchers against a worker acquiring jobs. Exactly one job per dedup id, and every returned id is a stored job.
  • dedup: concurrent dispatches do not fail while jobs are completed: dispatchers against workers completing jobs. Every added job is completed.
  • addDedupColumns should leave a dedup slot to its latest job, for Knex (SQLite, PostgreSQL) and Kysely (MySQL).
  • The Kysely MySQL group now runs the atomic dedup tests it skipped.

The two concurrency tests fail on main: Knex on PostgreSQL reports more added dispatches than stored jobs, and Kysely throws on PostgreSQL and MySQL.

The adapter specs (adapter.spec.ts, kysely_adapter.spec.ts, fake_adapter.spec.ts) pass on PostgreSQL, MySQL, SQLite and Redis.

The Knex adapter has no MySQL test group, so I checked both adapters on MySQL by hand, with a script running 6 dispatchers and 4 workers on 3 dedup ids. Dispatches and completions ended without an error, where main had 265 deadlocks.

Not in this PR

The same script shows popFrom hitting ER_LOCK_DEADLOCK on MySQL a few times per run with the Kysely adapter. It happens with plain jobs too, without any dedup, so it is a separate issue.

Commits

The first commit made the dispatch lock-free, with conditional updates. The second one brings the transaction and the row lock back, which is the design described above. They can be squashed.

🤖 Generated with Claude Code

A dispatch with .dedup() could enqueue a second job, or return the id
of a job that was never stored, when it raced with a worker acquiring
the job that owned the dedup id.

The unique index only covered pending and delayed jobs, so a job left
it once acquired. A concurrent dispatch that had read no owner then
either inserted a duplicate, or lost the insert and looked for the
winner among pending and delayed jobs only. Finding none, Knex returned
"added" with the id of the rolled back job and Kysely threw "Unable to
resolve concurrent dedup dispatch". MySQL has no partial index, so it
had no constraint at all, and its SELECT ... FOR UPDATE deadlocked with
the completion of the owner.

A plain unique index on (queue, dedup_id) now lets a single job own a
dedup id, whatever its status, on PostgreSQL, MySQL and SQLite: jobs
without dedup hold NULL, which is never a duplicate. Removing the job
releases the id.

The dispatch no longer opens a transaction. It reads the owner and
inserts when there is none; on a unique violation it reads the owner
again. The payload swap, the TTL refresh and the release of an expired
id are single UPDATE statements that carry the state they rely on in
their WHERE clause, so a job acquired in the meantime is not replaced
and an id extended in the meantime is not released. An expired id is
now released whatever the status of its job, retained ones included.
On MySQL, a dispatch picked as the deadlock victim among concurrent
inserts of the same id is retried.

addDedupColumns() creates the unique index on existing tables. Where
several jobs share a dedup id, the latest keeps it and the others lose
only their dedup id. The previous indexes are dropped. Without the
migration, sequential dispatches are still deduplicated, but
concurrent ones are not.
The dispatch read the owner of a dedup id without a lock, and guarded
each update with the state it had read. A change between the read and
the update could still be lost: a replace rewrites the whole data
column, so a job retried in that window got its attempts counter set
back.

The dispatch runs in a transaction again, and locks the owner with
SELECT ... FOR UPDATE before it skips, extends, replaces or releases
it. The updates are plain again.

The owner is found through the dedup index without a lock, then locked
through its primary key, as a worker locks a job. Locking it through
the dedup index takes the index entry then the row, while the removal
of a completed job takes the row then the index entry, and the two
deadlocked on MySQL.

FOR UPDATE cannot lock an owner that does not exist yet, so concurrent
first dispatches still race on the unique index. The one that loses
the insert runs its transaction again and locks the winner. The same
applies to a dispatch MySQL picks as a deadlock victim among
concurrent inserts of the same id.
@thetutlage

Copy link
Copy Markdown
Author

Why the owner is locked through its primary key (MySQL)

#resolveExistingDedup reads the owner of a dedup id in two steps instead of a single SELECT ... WHERE queue = ? AND dedup_id = ? FOR UPDATE. This only matters on MySQL. On PostgreSQL a row lock is the same lock, however the row was found.

In InnoDB, every index entry has its own lock, and a statement takes its locks in the order it walks the indexes:

Statement Locks first Locks second
SELECT ... WHERE queue = ? AND dedup_id = ? FOR UPDATE (dispatch) dedup index entry the row
DELETE ... WHERE id = ? AND queue = ? (completeJob, failJob) the row dedup index entry

Those are opposite orders. When a dispatch and the completion of the same job overlap, each holds one lock and waits for the other. InnoDB resolves the deadlock by aborting one of the two. In my runs it was completeJob, which leaves a finished job marked as active.

So the dispatch does this instead:

  1. Read the id of the owner through the dedup index, without a lock.
  2. Lock that row with WHERE id = ? AND queue = ? FOR UPDATE, which only takes the row lock.

The dispatch and the worker now both start with the row, so one waits for the other.

Step 1 holds no lock, so the row can change before step 2. That is why the code checks dedup_id again on the locked row. If the job was removed or released in the meantime, the dispatch falls through to the insert, and the unique index tells whether the id is free.

Measured with 6 dispatchers and 4 workers completing jobs on 3 dedup ids, Knex adapter on MySQL:

Version Deadlocks
main 262 on dispatch, 3 on completeJob
Unique index, owner locked through the dedup index 3 on completeJob
Unique index, owner locked through the primary key (this PR) 0

@thetutlage

Copy link
Copy Markdown
Author

Why the dispatch runs in a loop

FOR UPDATE can only lock a row that exists. On the first dispatch of a dedup id there is no owner yet, so two concurrent dispatches both read "no owner" and both insert. The unique index accepts one and rejects the other.

The rejected dispatch cannot carry on in the same transaction:

  • On PostgreSQL, a failed statement aborts the transaction. Every following statement fails until the rollback.
  • It still has to return the job id of the winner, and apply replace or extend to it when they were requested.

So it rolls back and runs the transaction again. The second pass finds the owner, locks it, and takes the regular "owner exists" path. The loop is only that: run the transaction again after losing the insert.

main handled the lost insert with a savepoint around it, followed by a separate lookup of the winner. That lookup filtered on pending and delayed jobs, which is where the wrong job id and the Unable to resolve concurrent dedup dispatch error came from. Running the transaction again reuses the code path that already handles an existing owner, instead of a second one written for the race.

The loop runs once, or twice when the insert is lost. A third pass needs another concurrent change at the same moment, such as the winner being completed and removed before it is read again.

Two more cases go through it:

  • MySQL deadlock victim. When several dispatches insert the same id right after its owner was removed, InnoDB can abort one of them with ER_LOCK_DEADLOCK. It runs again like a lost insert.
  • Duplicate job id. A unique violation can also come from the primary key. The dispatch would then loop forever, since no owner ever appears. After a violation, the code checks whether a job with this id exists, and rethrows the error if so.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant