---
title: "Replacing Cron Chaos with an Event-Driven Pipeline"
description: "Forty-one overnight cron jobs with no retries, no ownership and no visibility, moved onto EventBridge, SQS and Step Functions with dead-letter queues and idempotent consumers."
author: "Mohammad Abu Mattar"
canonical: https://mkabumattar.com/case-studies/post/cron-to-event-driven-pipeline
---

# Replacing Cron Chaos with an Event-Driven Pipeline

The handover note said "check the pipeline runs" and nothing else. Behind that sentence were 41 cron jobs spread across 6 EC2 hosts, some chained by `sleep 900` because someone had once needed job B to start after job A finished, and a crontab whose last editor had left the company. When something broke overnight, the first we heard was usually a customer asking why their report was empty.

This is what we replaced it with, and more usefully, what we got wrong on the way.

## The problem

The failure mode that finally forced the work was mundane. A vendor changed an SFTP filename, one job silently produced a zero-row file, six downstream jobs consumed that file happily, and by morning three internal dashboards were confidently reporting zero revenue. Nothing alerted, because nothing had technically failed. Every job exited 0.

Picking the system apart produced an uncomfortable inventory.

Retries did not exist. A job that failed at 02:14 was simply not done, and its output was missing until someone ran it by hand. Several jobs were not safe to run by hand twice, which nobody had documented, so re-running was a small act of faith.

Ordering was encoded in sleeps. Job A was scheduled at 02:00 and job B at 02:15, because A "usually takes ten minutes". When A grew to eighteen minutes, B read a half-written file. This had happened at least twice before and been fixed each time by moving B later.

Nothing was observable. Output went to `/var/log/cron` on whichever host had drawn the job, rotated weekly. There was no record of how long a job normally took, so there was no way to notice one taking four times as long.

Ownership was fiction. Of the 41 jobs, we could name a current owner for 12.

  The thing I remember is spending forty minutes just working out which host a
  job ran on. Not fixing it. Finding it.

The deepest problem was conceptual rather than operational. We eventually went through all 41 jobs and asked a single question: does this need to happen at a particular time, or does it need to happen when something else happens? **Three** were genuinely time-based. The other 38 were polling for a condition on a schedule, which is a workaround for not having an event.

## Constraints

No flag day. Finance closed the books off two of these jobs monthly, so "we are migrating the pipeline this weekend" was never an option. Every job had to move individually, with the old one still in place until the new one had run correctly through a full cycle.

We were on AWS and already used Lambda elsewhere, so managed services were preferred over anything we would have to operate. That ruled out introducing Airflow, which was the first suggestion in the room and would have replaced 41 unowned cron jobs with one unowned scheduler.

Two engineers, part time, alongside normal work. That budget shaped the sequencing more than anything technical: we needed each step to deliver something on its own, because there was a real chance of being pulled off it for a month.

The last constraint was one we set ourselves. Anything we built had to make the _next_ job easier to move, or it was not worth building. That is why the idempotency helper and the DLQ alarm came before most of the migrations.

## Architecture

The shape is unremarkable, which is the point. Events go on a bus, a queue absorbs bursts and retries, consumers are idempotent, and failure lands somewhere a human can see it.

_Event-driven pipeline on AWS_

EventBridge is the routing layer rather than the work layer. Producers publish one event shape and know nothing about consumers; rules decide who cares. That is what let us migrate incrementally, because adding a second consumer to an existing event is a rule, not a code change in the producer.

SQS sits between the rule and the worker on every path that does real work. This is the piece people skip when they wire EventBridge straight to Lambda, and it costs them the two things a queue gives you for free: a buffer when 4,000 events arrive at once, and somewhere for a failed message to wait instead of disappearing.

Step Functions handles anything with more than one step. The old pipeline encoded "A then B" as two cron entries fifteen minutes apart; a state machine encodes it as a dependency, retries each step independently, and shows you which step is currently failing without anyone reading a log.

Idempotency is the load-bearing decision, not the queue. SQS gives you at-least-once delivery, which means every consumer _will_ eventually see a message twice. If processing twice charges a card twice, no amount of retry tuning saves you. We put an idempotency key in DynamoDB with a conditional write, and after that a redelivery became boring.

The lifecycle a single message goes through is worth being explicit about, because most of our early bugs lived in it:

_One message, every path it can take_

## Implementation

1. **Inventory and classify, before touching anything.** List every job, its host, its schedule, its owner, and the one question that mattered: time-based or event-based? This produced the 3-versus-38 split that justified the whole project, and it took an afternoon.

2. **Build the shared pieces first.** An idempotency helper wrapping a DynamoDB conditional write, a Terraform module producing a queue plus its DLQ plus a depth alarm, and a structured log format. None of this migrated a job. All of it made every later migration a day instead of a week.

3. **Move one low-stakes job end to end.** We picked a job that regenerated a cache and could safely run twice. Getting one job fully through the new path surfaced the visibility-timeout problem and a permissions gap, cheaply.

4. **Run old and new in parallel, comparing output.** For each migrated job, both versions ran and wrote to different destinations, and a comparison job diffed them daily. Two jobs turned out to have been subtly wrong for months, which we only learned because the new implementation disagreed.

5. **Cut over per job, and delete the crontab entry the same day.** A commented-out cron entry is a job that comes back. If the new path had run clean through a full cycle, the old entry was removed, not disabled.

6. **Alarm on the DLQ, not on the job.** The old instinct was "alert when a job fails". The DLQ depth is a better signal: it fires when something has failed _and exhausted its retries_, which is the point at which a human is actually needed.

The idempotency wrapper is the piece worth showing, since it is what makes the rest safe:

```python

from botocore.exceptions import ClientError

_table = boto3.resource("dynamodb").Table(os.environ["IDEMPOTENCY_TABLE"])
TTL_SECONDS = 60 * 60 * 24 * 7  # a week is plenty; keys are only for dedupe

class AlreadyProcessed(Exception):
    """Raised when this key has already been completed."""

def claim(key: str) -> None:
    """Claim a key, or raise if someone already did the work.

    The conditional write is the whole mechanism: it is atomic, so two
    consumers racing on the same message cannot both win.
    """
    try:
        _table.put_item(
            Item={"pk": key, "expires_at": int(time.time()) + TTL_SECONDS},
            ConditionExpression="attribute_not_exists(pk)",
        )
    except ClientError as err:
        if err.response["Error"]["Code"] == "ConditionalCheckFailedException":
            raise AlreadyProcessed(key) from err
        raise

def handler(event, _context):
    for record in event["Records"]:
        # The key must come from the EVENT, not from anything generated at
        # runtime. A uuid4() here would make every redelivery look new.
        key = f"invoice-render:{record['messageAttributes']['invoiceId']['stringValue']}"
        try:
            claim(key)
        except AlreadyProcessed:
            continue  # delete the message; the work is already done
        render_invoice(record)
```

Derive the idempotency key from the event, never from anything created inside the handler. We got this wrong first time by including a timestamp, which made every retry look like new work and defeated the entire mechanism. The symptom was duplicate PDFs and a very confusing morning.

Queue configuration carries more of the reliability than the code does:

```hcl
resource "aws_sqs_queue" "invoice_render_dlq" {
  name                      = "invoice-render-dlq"
  message_retention_seconds = 1209600 # 14 days, the maximum
}

resource "aws_sqs_queue" "invoice_render" {
  name = "invoice-render"

  # Must exceed the handler's worst-case runtime, not its average. Ours p99s at
  # about 40s, so 180s leaves room for a slow dependency without the message
  # reappearing mid-flight and being processed twice.
  visibility_timeout_seconds = 180

  redrive_policy = jsonencode({
    deadLetterTargetArn = aws_sqs_queue.invoice_render_dlq.arn
    maxReceiveCount     = 5
  })
}

# The alarm that replaced "alert when a job fails".
resource "aws_cloudwatch_metric_alarm" "dlq_not_empty" {
  alarm_name          = "invoice-render-dlq-not-empty"
  namespace           = "AWS/SQS"
  metric_name         = "ApproximateNumberOfMessagesVisible"
  dimensions          = {QueueName = aws_sqs_queue.invoice_render_dlq.name}
  statistic           = "Maximum"
  period              = 300
  evaluation_periods  = 1
  threshold           = 0
  comparison_operator = "GreaterThanThreshold"
  alarm_actions       = [aws_sns_topic.oncall.arn]
}
```

- pipeline/
  - events/
    - **schema.json**
    - publish.py
  - shared/
    - idempotency.py
    - logging.py
  - consumers/
    - invoice_render/
    - ledger_export/
    - vendor_ingest/
  - terraform/
    - **queue/** (module: queue + DLQ + alarm)
    - eventbridge.tf
    - step_functions.tf
- runbooks/
  - dlq-redrive.md

## Results

Measured over the eight weeks before the work and the twelve weeks after, from CloudWatch and our incident log.

| Measure                               | Before                | After         | Source                         |
| ------------------------------------- | --------------------- | ------------- | ------------------------------ |
| Scheduled jobs                        | 41                    | 3             | measured                       |
| Jobs with automatic retry             | 0                     | all           | measured                       |
| Overnight pages                       | 9 in 8 weeks          | 2 in 12 weeks | measured, incident log         |
| Silent failures found after the fact  | 4                     | 0             | measured                       |
| Median time to identify a failing job | ~35 min               | under 2 min   | estimated, from incident notes |
| Jobs with a named owner               | 12 of 41              | all           | measured                       |
| Monthly compute cost                  | 6 always-on EC2 hosts | ~40% of that  | measured, Cost Explorer        |

The two remaining pages were both genuine: a vendor served malformed data for two days, and the DLQ caught every affected message. That is the system behaving correctly, and it is a different experience from being paged for something you then cannot locate.

The cost drop was a side effect rather than a goal. Six hosts sitting idle 23 hours a day became Lambda invocations plus one small Fargate task for the single job that outgrows the 15-minute limit.

Worth stating plainly: **eight jobs took longer to migrate than the other thirty combined.** Those eight had implicit ordering, shared mutable state in a scratch directory, or were not safe to run twice. The event model did not make them simple; it made their existing complexity visible and forced us to deal with it.

## Lessons

**"Scheduled" is usually a workaround for "no event exists yet."** Asking that one question of each job is what turned a vague modernisation project into a concrete list. It cost an afternoon and reframed everything. If you do nothing else from this write-up, do that inventory.

**Idempotency first, queues second.** We initially added SQS and enjoyed the retries, then spent a week chasing duplicate side effects. The queue is what makes retrying possible; idempotency is what makes retrying safe. In the wrong order, the queue is a duplicate-generating machine.

**The visibility timeout is a correctness setting, not a performance one.** Ours was 30 seconds against a handler that occasionally took 50. That produces two consumers working the same message, which looks exactly like a mysterious race condition in your own code. Set it from the p99, then add headroom.

**Alarm on the dead-letter queue, not on individual failures.** A transient failure that succeeds on retry is not an incident and should not wake anyone. A message that exhausted five attempts is. Moving the alarm to DLQ depth cut the overnight pages more than any reliability improvement did.

**Delete the old thing the same day.** We left three cron entries commented out "just in case". One got uncommented during an unrelated incident six weeks later and ran a job that had already been migrated, producing exactly the duplicate we had designed the idempotency layer to catch. It held, but the lesson is to remove rather than disable.

**Do not let a migration become a rewrite.** The temptation on every job was to fix its logic while moving it. We allowed this exactly twice, both times for jobs the comparison run proved were already wrong, and refused it everywhere else. The jobs we ported unchanged were the ones that went smoothly.

## Frequently Asked Questions

> **Why not just use Airflow or Dagster?**

Both are good tools, and both were the first suggestion. We passed because our problem was not orchestration complexity, it was that 38 of 41 jobs should not have been scheduled at all. Introducing a scheduler would have preserved that mistake in a nicer interface, and given us a new stateful service to operate and upgrade. If most of your jobs genuinely are time-based with real DAG dependencies, that calculus flips and a proper orchestrator is the right answer.

> **How do you handle a job that genuinely must run at a specific time?**

You keep it scheduled, just with the surrounding machinery. Our three survivors run on EventBridge Scheduler, which publishes an event onto the same bus everything else uses. That means they get the same queue, the same retries, the same DLQ and the same alarm as an event-driven consumer. The schedule is the trigger; nothing else about the path is special. This is much better than a crontab because the failure handling is shared rather than reimplemented per job.

> **What replaced the ordering that sleep timers used to provide?**

Explicit dependencies, in one of two forms. If step B needs step A's output, they belong in one Step Functions state machine where B is simply the next state. If B only needs to know that A happened, A publishes a completion event and B subscribes. The sleep-based version was really a guess about duration dressed up as a dependency, and it broke every time A got slower. Making the dependency explicit means it cannot silently degrade.

> **Is EventBridge plus SQS plus Lambda not more expensive than one cron host?**

Per invocation it is more expensive; overall it was cheaper for us, because six EC2 hosts ran continuously to do a few hours of nightly work. Our bill went to roughly 40% of what it was. The honest caveat is that this depends entirely on your duty cycle: a job that runs constantly is cheaper on a reserved instance than on Lambda, and a high-volume event stream can make SQS request charges noticeable. Model it against your own volumes rather than assuming serverless is cheaper.

> **How do you test an event-driven pipeline?**

Three layers, and the middle one matters most. Unit tests run the handler against a recorded event payload, which is easy and catches logic errors. Contract tests assert that a producer's event still matches the schema consumers expect, which is what stops one team's field rename breaking another team's consumer silently. Then a smoke test in a real environment publishes a synthetic event and asserts the expected side effect appeared. We skipped contract tests initially and paid for it once when a field became optional upstream.

> **What does the DLQ redrive process actually look like?**

It is a documented runbook, not a clever tool. Read a sample message from the DLQ to understand the failure, fix the cause, then redrive with SQS's built-in redrive or a small script that moves messages back onto the main queue in batches. The important detail is fixing the cause first: redriving into the same bug just refills the DLQ and wastes retries. We also cap redrive batch size, because dumping 40,000 parked messages back at once is its own outage.

> **Did you consider FIFO queues to guarantee ordering?**

We looked at them for two consumers and used them for neither. FIFO queues give you ordering within a message group at the cost of throughput and some operational sharp edges, and in both cases the real requirement was not "process in order" but "do not process the same entity concurrently". An idempotency key plus a per-entity lock addressed that without constraining the whole queue. If you have a genuine ordering requirement, for example applying a sequence of state transitions, FIFO is the right tool.

## References

- [Amazon EventBridge user guide](https://docs.aws.amazon.com/eventbridge/latest/userguide/eb-what-is.html)
- [EventBridge Scheduler](https://docs.aws.amazon.com/scheduler/latest/UserGuide/what-is-scheduler.html)
- [Amazon SQS visibility timeout](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-visibility-timeout.html)
- [Amazon SQS dead-letter queues](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-dead-letter-queues.html)
- [SQS dead-letter queue redrive](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-configure-dead-letter-queue-redrive.html)
- [AWS Step Functions error handling and retries](https://docs.aws.amazon.com/step-functions/latest/dg/concepts-error-handling.html)
- [DynamoDB condition expressions](https://docs.aws.amazon.com/amazondynamodb/latest/developerguide/Expressions.ConditionExpressions.html)
- [Making retries safe with idempotent APIs (AWS Builders' Library)](https://aws.amazon.com/builders-library/making-retries-safe-with-idempotent-APIs/)
- [Timeouts, retries and backoff with jitter (AWS Builders' Library)](https://aws.amazon.com/builders-library/timeouts-retries-and-backoff-with-jitter/)
- [AWS X-Ray for distributed tracing](https://docs.aws.amazon.com/xray/latest/devguide/aws-xray.html)
