Invisible watermarking sounds like a signal-processing problem. It is, until someone uploads a two-hour podcast and expects it back before their coffee cools. Then it becomes a distributed-systems problem. At Adori, I work on OrigID, a platform that embeds an inaudible ownership watermark into audio and video so a creator can later prove where a file came from. The encoder is CPU-heavy digital signal processing, and its cost grows linearly with the length of the file. On one worker, a long episode takes proportionally long. On AWS Lambda, a single invocation can't run longer than 15 minutes at all. This post walks through the pipeline we built to make processing time roughly independent of file length: split the file, watermark the pieces in parallel on Lambda, and stitch the result back together. The fan-out is the easy half. Most of this post is about the parts that aren't: where you're allowed to cut, how identity travels with the data, and how a fleet of stateless workers agrees that the job is finished. Design goals Latency that stays roughly flat as files get longer. No servers and no idle cost. Absorb bursts without capacity planning. Keep no media after processing. Intermediate files are deleted as soon as a job finishes. Every stage restartable, because events will be delivered more than once. Non-goals: live or streaming watermarking, GPUs, and exactly-once delivery. AWS doesn't offer exactly-once here, so we design for at-least-once instead. Watermarking, in one section Human hearing has blind spots. A loud sound masks quieter sounds near it in frequency and time, and psychoacoustic models predict how much energy can hide beneath that masking threshold at any moment. An audio watermark encoder spends that budget on a low-energy signal, added in the frequency domain, that carries an identifier: in our case, the owner's ID. The payload is protected with error correction and repeated continuously through the file, so a short clip, or a file that has been compressed, trimmed or converted to another format, still decodes. This builds on Adori's patented work on embedding inaudible data in audio (US 11,848,030 B2, Audio encoding for functional interactivity). I'll stay at that level. What matters for this post is what the encoder demands from the system around it: It works in fixed-size blocks. You can't cut the audio wherever you like. It wants uncompressed PCM at a known sample rate, so everything gets decoded first. Its cost is CPU-bound and roughly linear in duration, which is exactly the shape of work that parallelism pays off for. Each of those three shows up as a design decision below. Architecture at a glance Figure 1. One bucket with four prefixes, three kinds of Lambda functions, and a DynamoDB table for job state. The client asks a small signing API for an upload URL. It uploads the file directly to S3, under incoming/. The ObjectCreated event starts the splitter. The splitter creates a job record and writes N parts to parts/. Each part's ObjectCreated event starts its own encoder invocation. This is the fan-out. Encoders write to encoded-parts/ and report progress. Once every part is done, a merge step stitches them into output/ and cleans up. The user is notified. Meanwhile, the client polls the jobs table for progress. Why Lambda and S3 events rather than a queue and a fleet of workers on ECS or AWS Batch? Traffic is bursty: one large upload, then silence. With events there is nothing to scale and nothing to pay for between jobs. S3 notifications act as the work queue, and Lambda's concurrency acts as the scheduler. The costs are real: a 15-minute execution cap, a bounded /tmp disk, cold starts, and packaging native binaries. The 15-minute cap turns out to be useful, though, because it forces the split we wanted anyway. One practical rule: give every stage its own prefix, and never let a function write under the prefix that triggers it. Otherwise you've built an infinite loop with a credit card attached. Lambda now detects some recursive loops, but don't rely on it. Stage 1 — Ingest: upload straight to S3, identity in the envelope Media files are large, and there's no good reason to push them through your API. API Gateway caps request payloads at 10 MB and a synchronous Lambda invocation at 6 MB. Even without those limits, you'd be paying compute to shovel bytes you're about to write to S3 anyway. Instead, the signing API returns a presigned POST: a policy plus a signature that lets the browser upload one object directly to S3, under conditions we choose. Those conditions are the interesting part. A presigned POST can pin form fields, including user metadata, into the signed policy. If the client changes any of them, S3 rejects the upload. We use that to attach the job ID and the owner's ID to the object itself: job_id = uuid4().hex fields = { "x-amz-meta-job-id": job_id, "x-amz-meta-owner-id": str(owner_id), } conditions = [ {"x-amz-meta-job-id": job_id}, {"x-amz-meta-owner-id": str(owner_id)}, ["content-length-range", 1, MAX_UPLOAD_BYTES], ] post = s3.generate_presigned_post( Bucket=MEDIA_BUCKET, Key=f"incoming/{job_id}/{filename}", Fields=fields, Conditions=conditions, ExpiresIn=900, ) From here on, the identity travels with the data. The splitter reads it once and copies it onto every part it writes, and each encoder reads it from its own part's metadata. No worker ever asks a database "whose file is this?" on the hot path. Every worker has everything it needs in the object that triggered it. [Figure 2: upload images/fig2-identity-in-metadata.png] Figure 2. Owner identity is pinned by the upload signature, then carried as S3 object metadata through every stage. Stage 2 — Split: cut where the watermark lets you The splitter probes the file's duration first. Anything shorter than one part skips splitting and goes straight to a single encoder invocation. There's no reason to pay for the coordination described below. Longer files are decoded once to PCM at the canonical sample rate, and that one decision does a lot of work. Compressed formats can only be cut at their own frame boundaries (an MP3 frame is 1,152 samples), and they carry encoder delay and padding that turn into tiny gaps or clicks at every seam. PCM can be cut at any sample and joined back without loss. Decoding once also means every worker receives identical, sample-accurate input. Then comes the encoder's constraint: it works in fixed-size blocks. If a cut lands in the middle of a block, that block is split across two parts and neither side can carry it correctly. So the part length isn't "five minutes". It's the largest whole number of blocks that fits in the target duration: SR = 44_100 # canonical sample rate TARGET = 5 * 60 * SR # aim for ~5-minute parts, in samples part_len = (TARGET // BLOCK_LEN) * BLOCK_LEN # whole blocks only cuts = range(0, total_samples, part_len) # every cut lands on a block edge Figure 3. Cutting at fixed seconds splits a block across two parts. Cutting on block boundaries keeps every block intact, and the remainder becomes a shorter final part. Two details pay off later: Parts are self-describing. Each key encodes the job, the part's index and the last index, for example {jobId}_part_{i}_{last}.wav. Any worker knows where it sits and how many siblings it has without asking anyone, and the full list of parts can be rebuilt from any single key. Generate that list from the numbers rather than discovering it by pattern-matching a directory listing. The last part is the remainder, so it's almost always the shortest. That matters in the fan-in section. Stage 3 — Fan-out: let S3 events be the scheduler Writing N parts to parts/ produces N ObjectCreated events, and each one starts its own encoder invocation. There's no dispatcher and no queue for us to operate. S3 delivers the events, Lambda's internal asynchronous queue buffers them, and Lambda's concurrency decides how many run at once. Three settings matter: Reserved concurrency caps how many encoders run at once. It protects everything downstream, keeps the pipeline from starving other functions in the account, and bounds the cost of a runaway job. Events beyond the cap aren't lost: throttled asynchronous invocations keep being retried for up to six hours. Set the maximum event age and retry count deliberately rather than inheriting the defaults. Memory is CPU. Lambda allocates CPU in proportion to memory, reaching one full vCPU at 1,769 MB and six at 10,240 MB. A mostly single-threaded DSP binary stops getting faster once it has a full core, so benchmark the curve instead of maxing it out. /tmp is a budget. 16-bit stereo PCM at 44.1 kHz is 176,400 bytes per second, about 10.6 MB per minute. An encoder holding a 5-minute part plus its output needs roughly 110 MB. The splitter is the hungry one: the source file, the full PCM and all the parts together come to about 1.3 GB for a one-hour file, well past the 512 MB default. Ephemeral storage goes up to 10 GB. Size it for the longest file you accept, or stream instead of staging. The encoder and ffmpeg are native binaries, shipped in a Lambda layer built for Amazon Linux. Layers are unpacked under /opt, where the runtime finds them on the path. Layers are capped at 250 MB unzipped in total, so container images (up to 10 GB) are the escape hatch as the toolchain grows. Stage 4 — Fan-in: knowing when every part is done Fanning out is one S3 write per part. Fanning back in is the real design problem. N independent, stateless workers finish in an unpredictable order, and exactly one merge must run, only after all of them have succeeded. Nothing in a serverless system "waits" for anything, so you have to build the barrier yourself. There are four common ways to do it. [Figure 4: upload images/fig4-fan-in-patterns.png] Figure 4. Four fan-in patterns. The two on the left fail in subtle ways. The two on the right give you a single, reliable merge. A. A designated merger polls The worker that owns the last part index waits, polling S3 until every encoded part exists, then merges. It needs no extra infrastructure, which is why it's usually the first thing people build. But look at which worker got the job. The last part is the shortest (see Stage 2), so it usually finishes first, then sits there billing you while it waits for the longest parts. The barrier is bounded by that one worker's 15-minute timeout, and if it dies, nobody merges. B. Every worker checks the bucket After uploading, each worker lists encoded-parts/ and merges if the count equals the total. That fixes the question of who waits, but it introduces a race. Two workers finishing within milliseconds of each other can both see a complete set, and both merge. S3's strong read-after-write consistency doesn't save you here. It makes it more likely that both workers see the complete set. You need a single atomic decision, which is pattern C. C. An atomic counter in DynamoDB DynamoDB's UpdateItem is atomic per item, and ReturnValues="UPDATED_NEW" tells the caller what the value became. If each worker adds itself to a set of completed parts, exactly one worker observes the set becoming complete. Two refinements make this safe under at-least-once delivery. First, record which parts finished rather than how many. Adding the same part to a set twice changes nothing, while incrementing a counter twice silently completes the barrier early. Second, claim the merge with a conditional write, so a retried final worker can't start a second merge: table = dynamodb.Table(JOBS_TABLE) resp = table.update_item( Key={"jobId": job_id}, UpdateExpression="ADD doneParts :p", ExpressionAttributeValues={":p": {str(part_index)}}, # string set: idempotent ReturnValues="UPDATED_NEW", ) if len(resp["Attributes"]["doneParts"]) < total_parts: return # not the last part to finish try: table.update_item( Key={"jobId": job_id}, UpdateExpression="SET mergeStartedAt = :now", ConditionExpression="attribute_not_exists(mergeStartedAt)", ExpressionAttributeValues={":now": now_iso()}, ) except table.meta.client.exceptions.ConditionalCheckFailedException: return # another invocation already claimed the merge start_merge(job_id) There's no polling and no extra service, and the merge starts the moment the slowest part lands. D. An orchestrator: Step Functions Map Hand the barrier to AWS Step Functions. A Map state runs the encode step once per part, with its own concurrency limit, retry policy and error handling for each item. The next state, the merge, runs only when every iteration has finished. You also get an execution history to click through when something goes wrong. The price is more moving parts (S3 to EventBridge to a state machine) and, on Standard workflows, a charge for every state transition. Choosing A and B work in demos and fail in production, in ways that only show up under load or with unlucky timing. C is the smallest correct barrier: one table you probably already have, and about twenty lines of code. D is the right call when you want per-part retry policies and a visual execution history, or when the workflow grows more stages. Stage 5 — Merge and clean up Because every part is PCM, merging is a concatenation, not a re-encode: # parts.txt lists encoded parts in index order, e.g. file 'job_part_0_7.wav' ffmpeg -f concat -safe 0 -i parts.txt -c copy merged.wav ffmpeg -i merged.wav -c:a libmp3lame -b:a 128k output.mp3 The joins are sample-exact, with no generation loss and no gaps at the seams. The only lossy step is the final encode to the delivery format, and it happens exactly once. That final encode is also the first real test, since the watermark is designed to survive compression and this is the first place it has to. For video, only the audio track goes through this pipeline. The original video stream can be carried through untouched with -c:v copy. The merge then deletes everything under parts/ and encoded-parts/ for the job. Keeping no media after processing is a privacy promise, so it gets a second line of defence: an S3 lifecycle rule that expires those prefixes after a day catches any job that died halfway. Progress tracking without WebSockets The jobs table doubles as the progress bar. Once the splitter knows N, it writes the job's total number of steps: a fixed number for split and merge, plus a fixed number per part. Every stage then increments a completed-steps counter atomically, and progress is completed divided by total. Figure 5. The life of one job, from upload URL to download link. The client polls the jobs table until the job completes. The client polls a lightweight endpoint every few seconds: one GetItem per poll, no connection state, trivially cacheable. Push (WebSockets through API Gateway, or webhooks for API consumers) is worth it at larger scale, but polling a single item is hard to beat for simplicity. Keep status separate from percentage, though. Retries can over-count steps, so clamp the bar below 100% and treat the status field, not the arithmetic, as the source of truth for "done". A plain counter is fine here, unlike at the barrier: a duplicate increment only makes the bar jump, it never triggers a merge. Failure handling: at-least-once, everywhere S3 event notifications and asynchronous Lambda invocations are both at-least-once. Every stage has to assume it may run twice, occasionally at the same time as itself: Deterministic outputs. Part and output keys derive from the job ID and part index, so a rerun overwrites the same object instead of creating a sibling. Idempotent state changes. Use sets instead of counters wherever correctness depends on the count, and conditional writes wherever something must happen exactly once, like starting the merge or marking a job failed. An explicit retry policy. Configure the retry count and maximum event age for asynchronous invocations, and attach an on-failure destination or dead-letter queue, so an event that runs out of retries is visible rather than silently dropped. Fail fast on poison input. A corrupt or unsupported file fails the same way every time. Probe it up front, mark the job failed once, and tell the user, instead of burning retries. Watch for silence. In an event-driven pipeline the scariest failure isn't an error, it's nothing happening. A scheduled sweep that flags jobs stuck in progress beyond a reasonable bound catches what everything else misses. Timeout headroom. Size parts so the slowest one finishes comfortably inside Lambda's 15-minute limit, with margin for cold starts. The latency and cost model Let N be the number of parts, tᵢ the time to encode part i, and T_split and T_merge the split and merge times. On a single worker, processing takes roughly: T_serial ≈ t₁ + t₂ + … + t_N ≈ N × t̄ Fanned out, the parts run concurrently and you only wait for the slowest one: T_parallel ≈ T_split + max(tᵢ) + T_merge + cold starts Encoding time stops growing with file length. What's left grows slowly: split and merge are still proportional to duration, but they're I/O-bound decode, copy and concatenate operations, far cheaper per minute than DSP. For very long files they become the floor, which is Amdahl's law showing up on schedule. Cost is the pleasant surprise. Lambda bills per millisecond of execution at the configured memory, so N parts × t̄ costs about the same compute as one worker doing all of it. Parallelism adds only overhead: the split and merge, a few extra S3 and DynamoDB requests, and a cold start per concurrent worker. You're buying latency at almost no marginal cost, which is rare. Part size is the tuning knob. Smaller parts mean more parallelism but more per-part overhead and more seams. Larger parts mean fewer invocations but a slower slowest part and less headroom under the timeout. Pick it from measurements, and keep block alignment non-negotiable. Takeaways Let the algorithm choose the cut points. Your domain's invariants, here the block boundaries, decide where parallelism is allowed. Put identity in the envelope. Metadata pinned by the upload signature travels with the data and takes lookups off the hot path. Fan-out is a loop; fan-in is the design. Decide how exactly one worker learns that everyone is finished before you write the first encoder. Assume every event arrives twice. Deterministic keys, set-based completion tracking and conditional writes make that safe. In event-driven systems, silence is the failure mode. Alert on work that stops, not just on work that throws. On Lambda, parallelism costs about the same CPU-seconds. You're mostly buying latency. If you work on media pipelines, content provenance or large-scale serverless processing, I'd love to compare notes. Find me on LinkedIn.
Building a Parallel Distributed Audio Watermarking Pipeline on AWS Lambda
Full Article
Original Source
Read the full article at Hackernoon →KhanList aggregates and links to publicly available news content. We do not host full articles from third-party sources. Always verify important information with original sources.