skip to content

You attach a Lambda transformation to an Amazon Data Firehose delivery stream. What must the function return for each record, and what happens to records it cannot process?

level: middleimportance: nice to knowfreq 38%

answer

  1. one entry per input record
  2. echo the recordId back
  3. three results, three destinations
  4. Dropped is a filter, not a failure
  5. raw copy via source record backup

basics

~20 s

The function must return a records array with one entry per input record, each carrying the same recordId, a result of Ok, Dropped or ProcessingFailed, and base64-encoded data for Ok records. ProcessingFailed records are written to the delivery stream's S3 error output prefix.

solid answer

~60 s

Firehose invokes the function synchronously with a buffered batch, and the response contract is strict: return a `records` array with exactly one entry per input record, each echoing the **same `recordId`** Firehose supplied, plus a `result` and, for successful records, the transformed payload base64-encoded in `data`. The three `result` values mean different things: `Ok` delivers the record, `Dropped` filters it out deliberately and is not an error, and `ProcessingFailed` marks it as bad. If a `recordId` is missing or does not match, Firehose treats the whole invocation as failed. Failed invocations are retried (three times by default); records that still fail, and every record marked `ProcessingFailed`, are written to the delivery stream's S3 error output prefix rather than to the destination, so nothing is silently lost. Two practical limits shape the function: the synchronous invocation payload cap means the returned batch must stay under 6 MB, and enabling **source record backup** mirrors the raw untransformed records to a separate S3 location — the only way to recover from a bad transformation.

code

javascript · 23 lines
javascript
exports.handler = async (event) => {
  const records = event.records.map((record) => {
    const raw = Buffer.from(record.data, 'base64').toString('utf8');
    let payload;
    try {
      payload = JSON.parse(raw);
    } catch (err) {
      // unparseable: route to the S3 error output prefix, do not hide it
      return { recordId: record.recordId, result: 'ProcessingFailed' };
    }
    if (payload.eventType === 'healthcheck') {
      // intentional filtering: not an error
      return { recordId: record.recordId, result: 'Dropped' };
    }
    payload.ingestedAt = new Date().toISOString();
    return {
      recordId: record.recordId,
      result: 'Ok',
      data: Buffer.from(JSON.stringify(payload) + '\n').toString('base64'),
    };
  });
  return { records };
};

go deeper

for a junior

Know that Firehose can call a Lambda function on records before delivery, and that the function returns one entry per input record with a result value rather than a bare payload.

for a middle

Be able to state the response contract precisely — matching recordId, the three result values, base64-encoded data — and say where ProcessingFailed records end up.

for a senior

Show the operational view: alarm on the S3 error output prefix, enable source record backup before trusting a transformation, and keep the function cheap because it runs on every record in the feed.

for a principal

Own the policy question — how much logic belongs in an ingest-path transformation versus a downstream job, given that a bug there corrupts the only copy of the data and that the function's latency and cost scale with every event the platform emits.

## Where the transformation sits Amazon Data Firehose can call an AWS Lambda function on records *before* they are delivered. It is the extension point that turns Firehose from a dumb pipe into something that can normalise, enrich, filter or reshape a feed — flatten a nested payload, drop health-check noise, add a tenant id, or convert a non-JSON format into the JSON that record format conversion requires downstream. Crucially, the transformation runs on a **buffer**, not on single records. Firehose accumulates records, invokes the function synchronously with that batch, takes the response, and only then applies its delivery buffering and writes the result out. ## The response contract The event Lambda receives contains a `records` array; each element carries a `recordId` generated by Firehose and a base64-encoded `data` payload. The response must mirror it exactly: - one entry per input record — no more, no fewer; - the **same `recordId`** on each entry, so Firehose can match output to input; - a `result` field with one of three values; - for `Ok` records, a `data` field holding the transformed payload, base64-encoded. The three result values are not interchangeable: | `result` | Meaning | Where the record goes | |---|---|---| | `Ok` | Transformed successfully | Delivered to the destination | | `Dropped` | Intentionally filtered out | Nowhere — and it is not an error | | `ProcessingFailed` | Could not be processed | The S3 error output prefix | `Dropped` is the one people misuse. It exists so you can filter — discard debug events, sample a noisy feed — without polluting your error path. If you mark unparseable records `Dropped` you make real data loss invisible; mark them `ProcessingFailed` and they land in S3 where you can inspect and replay them. If the response omits a record, invents a `recordId`, or is malformed, Firehose cannot reconcile the batch and treats the entire invocation as failed. ## What happens on failure There are two failure layers. **Invocation failures** — the function errored, timed out, or returned something Firehose could not parse. Firehose retries the invocation, three times by default. If the retries are exhausted, the records in that batch are not dropped: they are written to the delivery stream's configured S3 error output location so the data survives. **Per-record failures** — records the function itself marked `ProcessingFailed`. These skip the destination and go to the same S3 error output location, wrapped with metadata describing why they were routed there. This is the design principle worth stating in an interview: *Firehose prefers to park bad data in S3 rather than lose it or block the pipeline*. Your job is to actually look at that prefix. An error path nobody monitors is the same as no error path, so alarm on objects appearing under it. ## Two limits that shape the function **Payload size.** Because the invocation is synchronous, the request and the returned batch are bound by Lambda's synchronous payload limit of 6 MB. A transformation that inflates records — expanding an abbreviated schema, adding verbose enrichment — can push a batch over the line even though the input fit. The lever is the transformation buffer hint, which is separate from the delivery buffering hints and is configured on the processor itself; lowering it makes each invocation smaller. **Timeout and cost.** The function runs on every record that passes through the stream, so its duration multiplies by volume. Keep it simple and side-effect free; if a transformation needs to call another service per record, that latency now sits inside your ingest path and the timeout becomes a delivery risk. ## Source record backup: the recovery you will wish you had When a transformation is enabled, Firehose offers **source record backup**, which writes the *untransformed* records to a separate S3 bucket or prefix in parallel with the transformed delivery. Enable it. Firehose stores nothing itself, so if the transformation logic was subtly wrong — dropped a field, mangled a timestamp — the transformed output in the destination is the only copy of that data, and it is wrong. Source record backup gives you the raw feed to reprocess from. ```javascript // minimal shape of the response Firehose expects return { records: [{ recordId: id, result: 'Ok', data: base64Payload }] }; ``` ## Common mistakes Generating new record ids; returning only the records that succeeded; forgetting to base64-encode the transformed payload; forgetting the trailing newline when the destination is S3 and a downstream reader expects newline-delimited JSON; and treating `Dropped` as a synonym for "failed". Each of these produces either a hard invocation failure or, worse, quiet data loss that nobody notices until a report looks wrong weeks later.

  • What is the difference between marking a record Dropped and marking it ProcessingFailed?
    `Dropped` means you deliberately filtered the record: it is not delivered and not treated as an error. `ProcessingFailed` means the record could not be processed, and Firehose writes it to the S3 error output prefix so you can inspect and replay it. Using `Dropped` for bad records turns real data loss into an invisible one.
  • The transformation is enabled and you later discover it mangled a field. How do you recover?
    Only if source record backup was enabled. Firehose retains nothing itself, so the transformed output is otherwise the sole copy. Source record backup mirrors the raw records to a separate S3 location, and you reprocess from there. Without it, the data is gone — which is why it should be on by default for any non-trivial transformation.
  • Why can a transformation that works on small batches start failing as traffic grows?
    The invocation is synchronous, so the batch sent to Lambda and the batch returned are bound by Lambda's 6 MB synchronous payload limit. A transformation that enriches records inflates the response, and a larger buffer can push it over. Lower the transformation buffer hint so each invocation carries less.

saying these in an interview costs you the question

  • Returns only the successfully transformed records
  • Generates fresh record ids instead of echoing them
  • Marks unparseable records Dropped, hiding data loss
  • Forgets to base64-encode the transformed payload
  • Assumes failed records are silently discarded by Firehose

context