Chunk, Embed, Index: Running RAG Ingestion as a Step Functions State Machine
Blog article/Blog archive

Chunk, Embed, Index: Running RAG Ingestion as a Step Functions State Machine

A re-run of a failed ingestion job left me with duplicate vectors and worse retrieval. Here's the Step Functions state machine I built instead: a Distributed Map over chunks, retries tuned to Bedrock throttling, content-hash idempotency, and a quarantine path for documents that just won't parse.

Aug 8, 202613 min read0 comments30 views
AWSStep FunctionsLambdaRAGBedrockTerraform

A customer uploaded a 412-page product manual to my app. The ingestion Lambda chewed through it for 14 minutes and 51 seconds, then died at the Lambda timeout. I looked at the logs, saw the timeout, shrugged, and hit re-run.

That was the mistake. The first run had already written about 900 chunks into the vector table before it died. The second run wrote them again, with fresh primary keys, because my chunk boundaries had shifted by a few characters after a whitespace change in the extractor. Nothing errored. The table just quietly grew.

I found out three days later, when a user asked a question and the top 5 retrieved chunks were the same paragraph five times. Retrieval had gotten worse, not better, because near-identical vectors crowded out everything else. Fixing it meant a DELETE against a document I couldn't fully identify, then a full re-ingest of that customer's library.

None of that was an embedding problem or a pgvector problem. It was an orchestration problem. My ingestion was a single Lambda doing a for-loop over chunks, with no notion of which chunks had already landed. So I rebuilt it as a Step Functions state machine, and the loop became a Distributed Map.

Ingestion is a different job than retrieval

I've written about the query side already in Building a Production RAG Pipeline on AWS Lambda and pgvector. That post is about what happens in the 300ms after a user hits enter. Retrieval is latency-bound, single-shot, and if it fails you just return an error and the user retries.

Ingestion has the opposite shape. It runs for minutes, not milliseconds. It calls an external embedding API hundreds or thousands of times per document. It writes state that persists forever. And when it fails halfway, the damage isn't a bad response, it's a corrupted index that poisons every future query.

That's a workflow, not a request handler. And Step Functions is the right tool for it, which is the same conclusion I reached in How I Use Step Functions to Orchestrate LLM Workflows Without Chaining Lambdas, except that post was about a fixed sequence of steps. Ingestion isn't a fixed sequence. The number of steps depends on how big the document is, and that changes everything about how you build it.

The stages, and why they're separate

The state machine has five stages:

  1. Extract. Pull text out of whatever landed in S3. PDF, DOCX, HTML, plain text.
  2. Plan. Split the text into chunks and write a manifest to S3. Return the manifest key, not the chunks.
  3. Index. A Distributed Map over the manifest. Each iteration embeds a batch of chunks and upserts them.
  4. Sweep. Delete rows from previous versions of this document that the current run didn't touch.
  5. Quarantine. The catch-all destination for documents that failed in a way retries won't fix.

Extract and Plan are separate on purpose. Extraction fails for boring reasons like a password-protected PDF or a scanned image with no text layer. Those are permanent failures, and I want them to fail before I've spent a cent on embeddings. Planning fails for basically no reason at all, so it gets one retry and that's it.

Extract Lambda Plan writes manifest to S3 Index — Distributed Map Batch A · ≤25 chunks Batch B · ≤25 chunks Batch C · MaxConcurrency 8 Sweep only if Map succeeded Quarantine SQS message + Fail state
Extract and Plan each run once. Index fans out across a Distributed Map sized to the Bedrock quota. Sweep only fires if every batch in the Map succeeded. Any Catch block, from any stage, routes straight to Quarantine.

Plan writes a manifest, not a payload

The first thing that bit me when I moved to Step Functions: state payloads cap at 256KB. A 412-page manual chunked at 500 tokens is roughly 1,100 chunks. That's a few megabytes of text. You cannot pass it between states.

So the Plan Lambda writes the chunk list to S3 and returns a pointer. The manifest is a JSON array, one object per chunk, which is exactly what the Distributed Map's ItemReader knows how to read.

import { createHash } from "node:crypto";
import { PutObjectCommand, S3Client } from "@aws-sdk/client-s3";

const hrr_s3 = new S3Client({});

export async function hrr_planChunks(event) {
  const { documentId, ingestRunId, text, metadata } = event;

  const chunks = hrr_chunkText(text, { maxTokens: 500, overlap: 50 });

  const manifest = chunks.map((content, index) => ({
    documentId,
    ingestRunId,
    chunkIndex: index,
    content,
    contentSha256: createHash("sha256").update(content).digest("hex"),
    metadata,
  }));

  const key = `manifests/${documentId}/${ingestRunId}.json`;

  await hrr_s3.send(
    new PutObjectCommand({
      Bucket: process.env.HRR_PIPELINE_BUCKET,
      Key: key,
      Body: JSON.stringify(manifest),
      ContentType: "application/json",
    })
  );

  return {
    documentId,
    ingestRunId,
    chunkCount: manifest.length,
    manifest: { bucket: process.env.HRR_PIPELINE_BUCKET, key },
  };
}

The contentSha256 field is the whole idempotency story, and I'll get to it. Compute it here, once, while the text is in front of you. Don't recompute it in the indexing step, because then a subtle difference in normalization gives you a different hash and the skip logic silently stops working.

Two things about the chunker itself. It has to be deterministic: same input text, same chunk boundaries, every time. And it has to be versioned. I keep a HRR_CHUNKER_VERSION constant and stamp it into the metadata, so when I change the splitting rules I can find every document that needs a full re-index instead of guessing.

Fan out with a Distributed Map, not a loop

Here's the state that replaced my for-loop. It's a Map in DISTRIBUTED mode, reading items straight out of the S3 manifest.

{
  "IndexChunks": {
    "Type": "Map",
    "ItemReader": {
      "Resource": "arn:aws:states:::s3:getObject",
      "ReaderConfig": { "InputType": "JSON" },
      "Parameters": {
        "Bucket.$": "$.manifest.bucket",
        "Key.$": "$.manifest.key"
      }
    },
    "ItemBatcher": {
      "MaxItemsPerBatch": 25,
      "BatchInput": {
        "documentId.$": "$.documentId",
        "ingestRunId.$": "$.ingestRunId"
      }
    },
    "MaxConcurrency": 8,
    "ToleratedFailurePercentage": 0,
    "ItemProcessor": {
      "ProcessorConfig": {
        "Mode": "DISTRIBUTED",
        "ExecutionType": "STANDARD"
      },
      "StartAt": "EmbedAndUpsertBatch",
      "States": {
        "EmbedAndUpsertBatch": {
          "Type": "Task",
          "Resource": "${hrr_embed_index_lambda_arn}",
          "TimeoutSeconds": 300,
          "Retry": [
            {
              "ErrorEquals": [
                "Bedrock.ThrottlingException",
                "Bedrock.ServiceQuotaExceededException",
                "Bedrock.ModelNotReadyException"
              ],
              "IntervalSeconds": 4,
              "MaxAttempts": 6,
              "BackoffRate": 2,
              "MaxDelaySeconds": 120,
              "JitterStrategy": "FULL"
            },
            {
              "ErrorEquals": [
                "Lambda.ServiceException",
                "Lambda.SdkClientException",
                "Lambda.TooManyRequestsException",
                "Postgres.ConnectionError"
              ],
              "IntervalSeconds": 2,
              "MaxAttempts": 4,
              "BackoffRate": 2,
              "JitterStrategy": "FULL"
            }
          ],
          "End": true
        }
      }
    },
    "ResultWriter": {
      "Resource": "arn:aws:states:::s3:putObject",
      "Parameters": {
        "Bucket": "${hrr_pipeline_bucket}",
        "Prefix": "map-results"
      }
    },
    "Next": "SweepStaleChunks"
  }
}

Four settings in there are doing most of the work.

MaxItemsPerBatch: 25 means each child execution gets 25 chunks, not one. Bedrock Titan embeddings take a single input per call, so the Lambda still loops internally, but the orchestration overhead drops by 25x. I'll show the cost difference below because it's larger than you'd guess.

MaxConcurrency: 8 is the throttle. This is the number I tune against my Bedrock quota, not against how fast I want the job to finish. With 25 chunks per batch and 8 concurrent batches, I'm asking for at most 200 embeddings in flight. Check your account's requests-per-minute quota for the embedding model in your region and work backwards from it. Setting this to 0 (unlimited) is how you discover what a throttling storm looks like.

ToleratedFailurePercentage: 0 means one failed batch fails the whole Map. That's deliberate. A partially indexed document is worse than an unindexed one, because retrieval will confidently return the half you did index. I'd rather quarantine the whole document and re-run it clean.

JitterStrategy: "FULL" matters more here than anywhere else in the pipeline. Without it, 8 concurrent batches all get throttled at the same moment, all back off by exactly 4 seconds, and all retry at exactly the same moment. Full jitter randomizes the delay across the window and spreads the retry storm out.

Retries that match the actual failure

The two Retry blocks exist because embedding throttles and Lambda hiccups need different treatment.

Throttling gets 6 attempts with a ceiling of 120 seconds. Bedrock quota pressure can last a while, especially if something else in your account is hammering the same model. A short backoff just burns your attempts during the exact window you should be waiting out.

Transient infrastructure errors get 4 attempts starting at 2 seconds. If a Lambda cold start or an RDS Proxy connection blip hasn't cleared in 30 seconds, it isn't transient.

What's deliberately not in either list: validation errors, unsupported file types, and text that exceeds the model's input limit. Retrying a 400 is just a slower way to fail. Those throw a typed error and get caught:

{
  "ExtractText": {
    "Type": "Task",
    "Resource": "${hrr_extract_lambda_arn}",
    "TimeoutSeconds": 300,
    "Retry": [
      {
        "ErrorEquals": ["Lambda.ServiceException", "Lambda.SdkClientException"],
        "IntervalSeconds": 2,
        "MaxAttempts": 3,
        "BackoffRate": 2
      }
    ],
    "Catch": [
      {
        "ErrorEquals": ["HrrUnsupportedDocument", "HrrEmptyDocument"],
        "Next": "QuarantineDocument",
        "ResultPath": "$.failure"
      },
      {
        "ErrorEquals": ["States.ALL"],
        "Next": "QuarantineDocument",
        "ResultPath": "$.failure"
      }
    ],
    "Next": "PlanChunks"
  }
}

Two Catch entries pointing at the same state looks redundant. It isn't. The first one matches on my own error names, so the quarantine record gets a clean reason like HrrUnsupportedDocument instead of States.TaskFailed. When you're looking at 200 quarantined documents on a Monday, that distinction is the difference between a two-minute triage and an hour of log spelunking.

To make that work, throw named errors from the Lambda:

class HrrUnsupportedDocument extends Error {
  constructor(message) {
    super(message);
    this.name = "HrrUnsupportedDocument";
  }
}

export async function hrr_extractText(event) {
  const text = await hrr_readDocument(event.bucket, event.key);

  if (text === null) {
    throw new HrrUnsupportedDocument(
      `No text layer in s3://${event.bucket}/${event.key}`
    );
  }
  if (text.trim().length < 20) {
    const err = new Error("Document produced fewer than 20 characters");
    err.name = "HrrEmptyDocument";
    throw err;
  }

  return { ...event, text, charCount: text.length };
}

The name property is what Step Functions matches on in ErrorEquals, not the class name. If you're bundling with esbuild and minifying, class names get mangled and your catch rules stop matching. Set name explicitly and it survives.

Idempotent indexing, so a re-run is free

This is the part that fixes the duplicate-vector mess I opened with. Three pieces: a stable key, a content hash, and a run ID.

CREATE TABLE hrr_doc_chunks (
  document_id      TEXT NOT NULL,
  chunk_index      INTEGER NOT NULL,
  content_sha256   TEXT NOT NULL,
  content          TEXT NOT NULL,
  embedding        vector(1024) NOT NULL,
  metadata         JSONB NOT NULL DEFAULT '{}',
  ingest_run_id    TEXT NOT NULL,
  updated_at       TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  PRIMARY KEY (document_id, chunk_index)
);

CREATE INDEX hrr_doc_chunks_run_idx
  ON hrr_doc_chunks (document_id, ingest_run_id);

The primary key is (document_id, chunk_index), so a chunk can only exist once. That alone would have saved me, since my duplicate rows all came from a BIGSERIAL key that happily accepted the same content twice.

The batch handler reads the existing hashes first, skips anything unchanged, and only pays Bedrock for chunks that actually moved:

export async function hrr_embedAndUpsertBatch(event) {
  const { BatchInput, Items } = event;
  const { documentId, ingestRunId } = BatchInput;

  const indexes = Items.map((item) => item.chunkIndex);

  const existing = await hrr_pool.query(
    `SELECT chunk_index, content_sha256
       FROM hrr_doc_chunks
      WHERE document_id = $1 AND chunk_index = ANY($2::int[])`,
    [documentId, indexes]
  );

  const seen = new Map(
    existing.rows.map((r) => [r.chunk_index, r.content_sha256])
  );

  let embedded = 0;
  let skipped = 0;

  for (const item of Items) {
    if (seen.get(item.chunkIndex) === item.contentSha256) {
      // Same text, same vector. Touch the run ID so the sweep spares it.
      await hrr_pool.query(
        `UPDATE hrr_doc_chunks
            SET ingest_run_id = $3, updated_at = NOW()
          WHERE document_id = $1 AND chunk_index = $2`,
        [documentId, item.chunkIndex, ingestRunId]
      );
      skipped += 1;
      continue;
    }

    const embedding = await hrr_embed(item.content);

    await hrr_pool.query(
      `INSERT INTO hrr_doc_chunks
         (document_id, chunk_index, content_sha256, content,
          embedding, metadata, ingest_run_id)
       VALUES ($1, $2, $3, $4, $5::vector, $6, $7)
       ON CONFLICT (document_id, chunk_index) DO UPDATE
         SET content_sha256 = EXCLUDED.content_sha256,
             content        = EXCLUDED.content,
             embedding      = EXCLUDED.embedding,
             metadata       = EXCLUDED.metadata,
             ingest_run_id  = EXCLUDED.ingest_run_id,
             updated_at     = NOW()`,
      [
        documentId,
        item.chunkIndex,
        item.contentSha256,
        item.content,
        JSON.stringify(embedding),
        item.metadata,
        ingestRunId,
      ]
    );
    embedded += 1;
  }

  return { documentId, embedded, skipped };
}

The skip path is why re-running a failed job is cheap now. When that 412-page manual times out at chunk 900, the retry embeds 200 chunks instead of 1,100. It also means a customer editing one paragraph in a long doc pays for one embedding, not eleven hundred.

Then the sweep. Every row the current run touched carries the current ingest_run_id. Anything left behind belongs to an older, longer version of the document:

DELETE FROM hrr_doc_chunks
 WHERE document_id = $1
   AND ingest_run_id <> $2;

One statement, and the stale tail is gone. This is the piece I didn't have before. My old code deleted by chunk_index >= chunkCount, which works only if the run completed. If it died at chunk 900, the delete never ran, and the leftovers stayed forever.

The sweep goes in its own state after the Map, because it must not run unless every batch succeeded. That's the other reason ToleratedFailurePercentage is 0. A partial Map success followed by a sweep would delete chunks that were never re-indexed.

Where broken documents go

Every Catch in the machine points at the same terminal state:

{
  "QuarantineDocument": {
    "Type": "Task",
    "Resource": "arn:aws:states:::sqs:sendMessage",
    "Parameters": {
      "QueueUrl": "${hrr_quarantine_queue_url}",
      "MessageBody": {
        "documentId.$": "$.documentId",
        "ingestRunId.$": "$.ingestRunId",
        "sourceKey.$": "$.key",
        "error.$": "$.failure.Error",
        "cause.$": "$.failure.Cause",
        "executionArn.$": "$$.Execution.Id",
        "failedAt.$": "$$.State.EnteredTime"
      }
    },
    "Next": "FailExecution"
  },
  "FailExecution": {
    "Type": "Fail",
    "Error": "HrrIngestionFailed",
    "Cause": "Document quarantined. See the quarantine queue for details."
  }
}

$$.Execution.Id is the context object, and it's the field that makes the queue useful. With the execution ARN on the message, I can open the exact failed run in the console from a queue item and see the input at every state.

Note the Fail state after the quarantine. It's tempting to end the execution successfully once you've recorded the failure, but then ExecutionsFailed stays at zero and your alarms never fire. Record the failure, then fail.

Once the underlying problem is fixed, Step Functions can redrive a failed execution from the state that broke, reusing the original input. The Extract and Plan work doesn't repeat, and thanks to the content hashes, neither does most of the embedding.

aws stepfunctions redrive-execution \
  --execution-arn arn:aws:states:us-east-1:111122223333:execution:hrr-rag-ingestion:abc123

That command only works on executions that are actually in a failed state, which is one more reason for the Fail at the end.

What the fan-out actually costs

Standard workflow state transitions run $0.025 per 1,000, so $0.000025 each. That sounds like nothing until you fan out per chunk.

For that 1,100-chunk manual:

  • One child execution per chunk: 1,100 children, roughly 3 transitions each, about 3,300 transitions. Call it $0.083 per document.
  • Batches of 25: 44 children, same 3 transitions each, about 132 transitions. Roughly $0.003 per document.
1 child / chunk $0.083 / doc Batches of 25 $0.003 / doc
Standard Step Functions transitions for the same 1,100-chunk document, one child execution per chunk versus batches of 25.

Twenty-eight times cheaper for a one-line config change. At a few hundred documents a month that's lunch money, but the ratio holds at any volume, and orchestration cost is the one line item people forget to model when they design a fan-out.

If your per-batch work is short and you don't need per-child execution history, switching ExecutionType to EXPRESS bills by duration and memory instead of transitions, which is cheaper again. I stay on STANDARD because when a batch fails I want the full history for those 25 chunks, and I'd rather pay $0.003 than guess.

The Terraform bits that are easy to miss

The state machine resource itself is ordinary. The IAM policy is where Distributed Map surprises people, because the state machine has to start executions of itself and read and write S3 on your behalf:

resource "aws_iam_role_policy" "hrr_ingestion_distributed_map" {
  name = "hrr-ingestion-distributed-map"
  role = aws_iam_role.hrr_ingestion.id

  policy = jsonencode({
    Version = "2012-10-17"
    Statement = [
      {
        Effect   = "Allow"
        Action   = ["states:StartExecution"]
        Resource = [aws_sfn_state_machine.hrr_ingestion.arn]
      },
      {
        Effect = "Allow"
        Action = ["states:DescribeExecution", "states:StopExecution"]
        Resource = [
          "${replace(aws_sfn_state_machine.hrr_ingestion.arn, "stateMachine", "execution")}:*"
        ]
      },
      {
        Effect   = "Allow"
        Action   = ["s3:GetObject", "s3:ListBucket"]
        Resource = [
          aws_s3_bucket.hrr_pipeline.arn,
          "${aws_s3_bucket.hrr_pipeline.arn}/*"
        ]
      },
      {
        Effect   = "Allow"
        Action   = ["s3:PutObject"]
        Resource = ["${aws_s3_bucket.hrr_pipeline.arn}/map-results/*"]
      }
    ]
  })
}

That self-referencing states:StartExecution is the one people miss. Without it the Map fails immediately with an access denied that names your own state machine, which reads like a bug in AWS until you realize a Distributed Map literally runs its iterations as child executions.

The trigger is an S3 event notification straight to EventBridge, then an EventBridge rule that starts the machine. No glue Lambda in between:

resource "aws_cloudwatch_event_rule" "hrr_document_uploaded" {
  name = "hrr-document-uploaded"

  event_pattern = jsonencode({
    source        = ["aws.s3"]
    "detail-type" = ["Object Created"]
    detail = {
      bucket = { name = [aws_s3_bucket.hrr_uploads.id] }
      object = { key = [{ prefix = "incoming/" }] }
    }
  })
}

resource "aws_cloudwatch_event_target" "hrr_start_ingestion" {
  rule     = aws_cloudwatch_event_rule.hrr_document_uploaded.name
  arn      = aws_sfn_state_machine.hrr_ingestion.arn
  role_arn = aws_iam_role.hrr_events_invoke_sfn.arn

  input_transformer {
    input_paths = {
      bucket = "$.detail.bucket.name"
      key    = "$.detail.object.key"
      etag   = "$.detail.object.etag"
    }

    input_template = <<EOT
{
  "bucket": <bucket>,
  "key": <key>,
  "documentId": <key>,
  "ingestRunId": <etag>
}
EOT
  }
}

Using the object ETag as the ingestRunId is a small trick that pays off. The ETag changes when the file content changes, so two uploads of the identical file produce the identical run ID. Combined with the content hashes, a duplicate upload embeds nothing and sweeps nothing. It's a no-op that costs a few state transitions.

Three things I got wrong first

I passed chunks through state. Worked fine on the test fixtures, which were all under 40KB. The first real PDF hit States.DataLimitExceeded and the error message points at the state, not at the payload size, so it took me longer than I want to admit. Manifest in S3 from the start.

I set MaxConcurrency to 40. The job finished fast for about a minute, then every batch started throttling at once, retries stacked, and the whole Map ground slower than the sequential version had been. Eight concurrent batches, sized to my quota, finishes sooner in wall time than forty fighting each other.

I forgot Distributed Map writes execution history to S3. The map-results prefix accumulates one manifest and result set per run. After a few thousand documents that's a lot of small objects nobody will ever read. A 30-day lifecycle rule on that prefix costs one Terraform block and prevents a slow, boring bill.

The ingestion pipeline has been running on this shape for a while now. Documents that fail land in a queue with the reason attached instead of vanishing. Re-runs are cheap because unchanged chunks skip the embedding call. And the duplicate vector problem is structurally impossible, because (document_id, chunk_index) is a primary key and a run ID sweeps whatever's left over.

Hope you enjoyed this one. If you're building RAG ingestion on AWS and you've found a better way to handle the fan-out, or you just want to argue about batch sizes, come find me on X at https://x.com/harundotdev.

Get the next note

I email when a new post goes up. One send a week, and only if there's something new.

Want this applied on your account? Start with Infrastructure as Code.

Related reading

More posts, sliding underneath the article

Kept below the post instead of in a sidebar, with a slow continuous motion for a cleaner editorial feel.

On this post

Comments

A reply stays under the note it answers.

0 comments

No comments yet.

If you have a note on Chunk, Embed, Index: Running RAG Ingestion as a Step Functions State Machine, sign in and leave it.