When you're building a Shopify app, webhooks are your lifeline. Orders, products, customers, checkouts — Shopify sends you a webhook for everything. But here's the problem: Shopify expects a 200 response within 5 seconds, or it retries. And if you process webhooks synchronously, a single slow handler blocks your entire server.
Our solution? A custom async job queue built entirely on PostgreSQL. No Redis. No Bull. No external dependencies. Just a database table, some clever SQL, and a lot of edge case handling.
This is the story of how we built it, why we made these choices, and the production bugs that taught us the hard way.
The Problem: Shopify Webhooks Are Brutal
Shopify's webhook contract is simple but unforgiving:
- Shopify sends a POST request to your endpoint
- You must respond with 200 OK within 5 seconds
- If you don't respond, Shopify retries (up to 19 times over 48 hours)
- If you respond with anything other than 200, Shopify retries
- Shopify doesn't guarantee delivery order
- Shopify can send the same webhook multiple times (at-least-once delivery)
Our discount app, DealCraft, needs to do a lot of work when an order webhook arrives:
- Match the order to discount rules
- Record analytics
- Update customer usage counts
- Generate post-purchase discount codes
- Trigger loyalty point awards
- Check referral rewards
- Sync metafields back to Shopify
This takes 2-10 seconds per order. If we processed synchronously, we'd timeout on every webhook. So we needed a queue.
Why Not Redis?
The obvious choice is Redis + Bull (or BullMQ). It's the standard Node.js job queue. But we had three constraints:
1. Single-server deployment
DealCraft runs on a single VPS (we're bootstrapping, not VC-funded). Adding Redis means another service to manage, monitor, and backup. We already have PostgreSQL — why add complexity?
2. Durability matters more than speed
Redis is in-memory. If the server crashes, you lose unprocessed jobs (unless you configure AOF persistence, which adds latency). For webhook processing, losing a job means losing an order's discount attribution. That's unacceptable.
3. We need idempotency at the database level
Shopify can send the same webhook multiple times. We need to guarantee that processing the same webhook twice doesn't create duplicate discount codes or double-count analytics. Redis queues don't give you this for free — you still need to implement deduplication logic.
PostgreSQL gives us all three: it's already running, it's durable by default, and we can enforce idempotency with unique constraints.
The Schema: Idempotency by Design
Here's the WebhookJob table:
model WebhookJob {
id String @id @default(uuid())
topic String @map("topic")
shopId String @map("shop_id")
shopDomain String @map("shop_domain")
webhookId String? @map("webhook_id")
payload Json @map("payload")
status String @default("pending") // pending | processing | completed | failed
attempts Int @default(0)
lastError String? @map("last_error")
createdAt DateTime @default(now()) @map("created_at")
updatedAt DateTime @updatedAt @map("updated_at")
processedAt DateTime? @map("processed_at")
retryAfter DateTime? @map("retry_after")
@@index([status, createdAt])
@@index([status, updatedAt])
@@unique([topic, shopId, webhookId]) // The magic line
@@map("webhook_jobs")
}
That @@unique([topic, shopId, webhookId]) constraint is the key. Shopify sends an X-Shopify-Webhook-Id header with every webhook — a unique identifier for that specific event. By making [topic, shopId, webhookId] unique, the database itself prevents duplicate processing.
One caveat: webhookId is nullable (String?). In PostgreSQL, NULL != NULL in unique constraints, so rows with webhookId = NULL bypass deduplication. In practice, Shopify always sends the header for standard webhooks (ORDERS_CREATE, PRODUCTS_CREATE, etc.), so this only affects internally-generated events (like our custom FLOW:* topics). If you rely on this pattern, either make the field non-nullable or generate synthetic IDs for events that lack one.
If Shopify retries the same webhook (or sends it twice due to a network glitch), the second INSERT throws a P2002 unique constraint violation. We catch it and silently ignore it. No duplicate jobs. No complex deduplication logic in application code.
This is idempotency by schema design, not by application logic.
Security: HMAC Verification at the Edge
Before any of this queue logic runs, we need to verify the webhook actually came from Shopify. Every webhook endpoint in our app goes through authenticate.webhook(request) from the Shopify SDK, which validates the X-Shopify-Hmac-Sha256 header against the request body using our API secret. If the signature doesn't match, we return 401 before the payload ever touches our queue.
This is handled by the SDK so we don't have to think about it in every route — but it's worth understanding that the security boundary is at the HTTP edge, not inside the queue.
Enqueueing: Fire and Forget (Almost)
When a webhook arrives at our endpoint, we do this:
async function handleWebhook(topic: string, shopDomain: string, payload: any, webhookId: string) {
const shop = await findShopByDomain(shopDomain);
if (!shop) return; // Unknown shop, ignore
try {
// Persist the job to the database
await prisma.webhookJob.create({
data: {
topic,
shopId: shop.id,
shopDomain,
webhookId,
payload,
status: "pending",
},
});
// Trigger immediate consumption (non-blocking)
setImmediate(() => consumeWebhookBatch(20));
} catch (error) {
if (error.code === "P2002") {
// Duplicate webhook — already enqueued, ignore
return;
}
// Database error — fall back to synchronous processing
await dispatch(topic, shopDomain, payload);
}
}
Three things to notice:
We persist first, then consume. The
setImmediateis just an optimization — it triggers consumption immediately so the job doesn't wait for the next cron tick. But ifsetImmediatefails, the job is still in the database and will be picked up by the cron fallback. One caveat: if 50 webhooks arrive in rapid succession (bulk operations, flash sales), we fire 50 concurrentconsumeWebhookBatchcalls. CAS prevents duplicate processing, but each call still runsrecoverWebhookJobs()and afindMany— wasted DB queries. In production, this hasn't been a problem at our scale (~100 jobs/min), but a debounce or advisory lock would be prudent at higher throughput.We catch P2002 specifically. If the unique constraint fires, it means this webhook was already enqueued. We return silently — no error, no retry.
We have a fallback — with a narrow duplicate window. If the database is down (rare, but it happens), we fall back to synchronous
dispatch(). This blocks the webhook response, but at least the business logic runs. Better to timeout than to lose the job entirely. The catch block doesn't just fire on "DB is down" — Prisma throws on connection pool exhaustion, network timeouts, and other errors. Some of these (like connection timeouts) may mean the INSERT actually succeeded but the response didn't reach us. In that case, the sync fallback processes the event and the persisted job will be consumed later — a duplicate. Our processors are idempotent by design (checking existing records before creating), so this window is tolerable. But it's worth understanding: the sync fallback is a last resort, not a safe path. It bypasses all queue protections (CAS, heartbeat, lease) and runs the full processing chain synchronously.
After enqueueing, we return 200 to Shopify immediately. Total time: ~50ms (one database INSERT). Shopify is happy.
CAS Claiming: The Heart of the Queue
Here's where it gets interesting. Multiple consumers might be running (we have a cron job that runs every minute, plus the setImmediate trigger). How do we prevent two consumers from processing the same job?
Compare-And-Swap (CAS) claiming.
async function consumeWebhookBatch(limit: number = 20) {
// 1. Recover any crashed jobs (expired leases)
await recoverWebhookJobs();
// 2-3. CAS claim inside a transaction
const claimedJobs = await prisma.$transaction(async (tx) => {
// 2. Fetch pending jobs
const pendingJobs = await tx.webhookJob.findMany({
where: {
status: "pending",
attempts: { lt: MAX_ATTEMPTS },
},
orderBy: { createdAt: "asc" },
take: limit,
});
const claimed: WebhookJob[] = [];
// 3. Atomic claim per job
for (const job of pendingJobs) {
const result = await tx.webhookJob.updateMany({
where: {
id: job.id,
status: "pending",
attempts: job.attempts, // The CAS token
},
data: {
status: "processing",
attempts: { increment: 1 },
updatedAt: new Date(),
},
});
if (result.count === 1) {
// We successfully claimed this job
claimed.push({ ...job, attempts: job.attempts + 1 });
}
// If count === 0, another consumer already claimed it — skip
}
return claimed;
});
// 4. Process claimed jobs
for (const job of claimedJobs) {
await processJob(job);
}
}
The magic is in the updateMany WHERE clause: { id, status: "pending", attempts: job.attempts }.
-
idensures we're updating the right job -
status: "pending"ensures it hasn't been claimed yet -
attempts: job.attemptsis the CAS token — it's the generation number
When we read the job, attempts is, say, 0. We try to update it with attempts: 0 in the WHERE clause. If another consumer has already claimed it (incrementing attempts to 1), the WHERE clause won't match, and updateMany returns count: 0. We skip the job.
This is optimistic concurrency control — we assume the job is available, try to claim it, and back off if we're wrong. The SELECT and UPDATEs run inside a $transaction, which ensures they share the same database connection. But the real safety net is PostgreSQL's row-level locking: even without the transaction, each updateMany is individually atomic, and the CAS condition ensures only one consumer wins.
Lease Protection: What If the Consumer Crashes?
CAS claiming prevents duplicate processing, but what if a consumer crashes mid-processing? The job is stuck in status: "processing" forever.
Solution: lease-based expiration with heartbeat renewal.
const WEBHOOK_LEASE_MS = 5 * 60_000; // 5 minutes
const WEBHOOK_HEARTBEAT_MS = 60_000; // 1 minute
// After claiming a batch, start a heartbeat
const heartbeatTimer = setInterval(() => {
renewHeartbeat(claimedJobs);
}, WEBHOOK_HEARTBEAT_MS);
heartbeatTimer.unref(); // Don't prevent process exit
async function renewHeartbeat(jobs: WebhookJob[]) {
const owners = jobs.map(j => ({ id: j.id, attempts: j.attempts }));
await prisma.webhookJob.updateMany({
where: {
status: "processing",
updatedAt: { gt: new Date(Date.now() - WEBHOOK_LEASE_MS) },
OR: owners,
},
data: {
updatedAt: new Date(),
},
});
}
Every 60 seconds, we touch the updatedAt field on all outstanding jobs. This renews their lease.
If a consumer crashes, the heartbeat stops. After 5 minutes, the job's lease expires (updatedAt <= now - 5min). The next consumer's recoverWebhookJobs() call finds it and resets it to pending:
const MAX_ATTEMPTS = 3;
const RECOVER_LIMIT = 50;
async function recoverWebhookJobs() {
const expiredWhere = {
status: "processing",
updatedAt: { lte: new Date(Date.now() - WEBHOOK_LEASE_MS) },
};
// Batch limit: recover at most 50 at a time to avoid overwhelming the consumer
const expiredJobs = await prisma.webhookJob.findMany({
where: expiredWhere,
select: { id: true, attempts: true },
take: RECOVER_LIMIT,
});
const failIds = expiredJobs
.filter(j => j.attempts >= MAX_ATTEMPTS)
.map(j => j.id);
const requeueIds = expiredJobs
.filter(j => j.attempts < MAX_ATTEMPTS)
.map(j => j.id);
if (failIds.length > 0) {
await prisma.webhookJob.updateMany({
where: { id: { in: failIds }, status: "processing" },
data: {
status: "failed",
lastError: "Processing lease expired after maximum attempts",
processedAt: new Date(),
},
});
}
if (requeueIds.length > 0) {
await prisma.webhookJob.updateMany({
where: { id: { in: requeueIds }, status: "processing" },
data: {
status: "pending",
lastError: "Processing lease expired; retry pending",
retryAfter: null, // Clear any backoff timer
},
});
}
}
This is eventual consistency for job processing. If a consumer crashes, the job is automatically requeued. No manual intervention.
The Pre-Dispatch Lease Check
There's one more edge case: what if the lease expires during processing, right before we try to update the job status?
Example:
- Consumer A claims job (attempts: 1, status: processing)
- Consumer A starts processing (takes 6 minutes due to slow Shopify API)
- Consumer A's heartbeat hasn't fired yet (or the server is under load)
- Consumer B's
recoverWebhookJobs()sees the expired lease and resets the job topending - Consumer B claims the job (attempts: 2, status: processing)
- Consumer A finishes processing and tries to mark the job as
completed
Now we have a race condition. Consumer A might overwrite Consumer B's work.
Solution: check the lease before every status transition.
async function processJob(job: WebhookJob) {
const owner = { id: job.id, status: "processing", attempts: job.attempts };
const activeOwner = () => ({
...owner,
updatedAt: { gt: new Date(Date.now() - WEBHOOK_LEASE_MS) },
});
try {
// Check lease before dispatching
const leaseCheck = await prisma.webhookJob.updateMany({
where: activeOwner(),
data: { updatedAt: new Date() }, // Renew lease
});
if (leaseCheck.count !== 1) {
// Lease lost — another consumer took over
return;
}
// Dispatch to the appropriate handler
await dispatch(job.topic, job.shopDomain, job.payload);
// Mark as completed (with CAS + lease check)
await prisma.webhookJob.updateMany({
where: activeOwner(),
data: {
status: "completed",
processedAt: new Date(),
},
});
} catch (error) {
// Handle failure (with CAS + lease check)
await handleJobFailure(job, error);
}
}
Every terminal state transition (completed, failed) includes the CAS token (attempts) and the lease check (updatedAt > now - 5min). If either fails, we back off. This prevents zombie consumers from corrupting state. The heartbeat and pre-dispatch check use lighter guards — they only need to confirm the lease hasn't expired, not enforce full CAS.
Retry Logic: Three Strikes with Exponential Backoff
When a job fails, we reset it to pending — but not for immediate retry. Hammering a failing API three times in rapid succession doesn't help anyone. Instead, we use exponential backoff via a retryAfter timestamp:
const BACKOFF_BASE_MS = 30_000; // 30 seconds
// activeOwner() builds the WHERE clause for terminal state transitions:
// it combines the CAS token (attempts) with a lease check (updatedAt).
const owner = { id: job.id, status: "processing", attempts: job.attempts };
const activeOwner = () => ({
...owner,
updatedAt: { gt: new Date(Date.now() - WEBHOOK_LEASE_MS) },
});
// In the failure handler:
const nextStatus = job.attempts >= MAX_ATTEMPTS ? "failed" : "pending";
const retryAfter = nextStatus === "pending"
? new Date(Date.now() + BACKOFF_BASE_MS * Math.pow(2, job.attempts - 1))
: null;
await prisma.webhookJob.updateMany({
where: activeOwner(),
data: { status: nextStatus, lastError: errMsg, retryAfter },
});
The consumer skips jobs whose retryAfter is in the future:
const pending = await tx.webhookJob.findMany({
where: {
status: "pending",
attempts: { lt: MAX_ATTEMPTS },
OR: [
{ retryAfter: null },
{ retryAfter: { lte: new Date() } },
],
},
// ...
});
This gives us: 1st failure → 30s delay, 2nd failure → 60s delay, 3rd failure → dead-letter. Most transient errors (Shopify rate limits, temporary network issues) resolve within the first backoff window.
We use MAX_ATTEMPTS = 3. After 3 failures, the job is marked failed (dead-letter). An admin can inspect the lastError field and manually retry or investigate.
Why 3 attempts? Because most webhook failures are transient: Shopify API timeouts, network glitches, temporary rate limits. If it fails 3 times, it's probably a bug or a data issue that needs manual intervention.
The Dispatch Router
Once a job is claimed, we route it to the appropriate handler based on the topic:
async function dispatch(topic: string, shopDomain: string, payload: any) {
switch (topic) {
case "ORDERS_CREATE":
await handleOrderCreate(shopDomain, payload);
break;
case "ORDERS_PAID":
await handleOrderPaid(shopDomain, payload);
break;
case "PRODUCTS_CREATE":
await handleProductCreate(shopDomain, payload);
break;
case "APP_SUBSCRIPTIONS_UPDATE":
await billingService.handleSubscriptionUpdate(shopDomain, payload);
break;
case "FLOW:ABANDONED_CART":
await handleAbandonedCart(shopDomain, payload);
break;
// ... 10+ more topics
default:
throw new Error(`Unknown webhook topic: ${topic}`);
}
}
The ORDERS_CREATE handler is the most complex — it does discount attribution, analytics, loyalty points, post-purchase codes, and more. It takes 2-10 seconds. This is why we needed the queue in the first place.
Production Bugs That Taught Us Lessons
Bug 1: The Heartbeat That Didn't Unref
Symptom: The Node.js process wouldn't exit after processing webhooks. It hung forever.
Cause: The heartbeat setInterval was keeping the event loop alive. We forgot to call .unref() on the timer.
Fix:
const heartbeatTimer = setInterval(() => renewHeartbeat(jobs), WEBHOOK_HEARTBEAT_MS);
heartbeatTimer.unref(); // Allow process to exit
Bug 2: The Duplicate Post-Purchase Codes
Symptom: Some orders got two post-purchase discount codes instead of one.
Cause: The handleOrderCreate processor wasn't idempotent. If the job was retried (due to a transient error), it would create a second discount code.
Fix: We added an idempotency check — query existing RuleLog entries for the order, and skip already-logged rules:
const existingLogs = await prisma.ruleLog.findMany({
where: { orderId: order.id },
});
const logsToCreate = allLogs.filter(
log => !existingLogs.some(existing => existing.ruleId === log.ruleId)
);
if (logsToCreate.length > 0) {
await prisma.ruleLog.createMany({
data: logsToCreate,
skipDuplicates: true, // Belt and suspenders
});
}
Bug 3: The Zombie Consumer
Symptom: Jobs were stuck in processing even though the consumer had crashed.
Cause: The consumer crashed before the heartbeat could renew the lease, and the recoverWebhookJobs() function wasn't being called frequently enough.
Fix: We call recoverWebhookJobs() at the start of every consumeWebhookBatch(), not just in the cron job. This ensures that crashed jobs are recovered immediately.
Performance: Is PostgreSQL Fast Enough?
Yes. Here are the numbers from our production deployment:
- Enqueue latency: ~20ms (one INSERT)
- Claim latency: ~50ms (one SELECT + N UPDATEs, where N ≤ 20)
- Throughput: ~100 jobs/minute (limited by processing time, not queue overhead)
- Table size: ~10,000 rows (we cleanup jobs older than 7 days)
The bottleneck is the processing logic (Shopify API calls), not the queue itself. PostgreSQL adds maybe 50ms of overhead per job — negligible compared to the 2-10 seconds of processing time.
When Would You Use Redis Instead?
Our PostgreSQL queue works great for a single-server deployment processing ~100 webhooks/minute. But there are scenarios where Redis would be better:
Multi-server deployment — If you have multiple app servers, Redis provides a centralized queue. PostgreSQL works, but you need to coordinate consumers across servers (which we do with CAS, but it's more complex).
High throughput — If you're processing 10,000+ jobs/second, Redis is faster (in-memory vs. disk).
Advanced features — Redis queues (Bull, BullMQ) give you priority queues, delayed jobs, rate limiting, and more out of the box.
For our use case (single server, ~100 jobs/minute, need durability and idempotency), PostgreSQL is the right choice.
The Takeaways
Idempotency by schema design — Use unique constraints to prevent duplicate processing, not application logic. Know the NULL caveat.
CAS claiming — Use a generation token (
attempts) to prevent duplicate consumption. EachupdateManyis individually atomic via PostgreSQL row-level locking; wrapping in$transactionshares the connection but isn't the safety net — the CAS condition is.Lease-based expiration — Use
updatedAtas a lease timestamp, with heartbeat renewal and automatic recovery. Batch-limit the recovery and guardupdateManywithstatus: "processing"to prevent racing a concurrent completion.Exponential backoff — Don't retry failures immediately. A
retryAftertimestamp gives transient errors time to resolve without wasting attempts.Fallback to synchronous — If the queue is unavailable, process synchronously as a last resort. Accept that the duplicate window is wider than the normal path (you can't distinguish "INSERT failed" from "INSERT succeeded but response timed out"). Make processors idempotent.
Beware implicit fan-out —
setImmediateper webhook is fine at low volume, but creates concurrent consumer storms under burst traffic. Debounce or use an advisory lock if throughput grows.PostgreSQL is good enough — For most use cases, you don't need Redis. PostgreSQL gives you durability, idempotency, and sufficient performance.
This post is based on our production experience building a Shopify app with TypeScript and Rust. If you're designing a webhook queue yourself, or have different takes on idempotency, CAS claiming, lease recovery, or exponential backoff, I'd love to hear your perspective in the comments.
Top comments (0)