Skip to content

feat(sdk): add distributed map operation - #1

Draft
nvasiu wants to merge 2 commits into
mainfrom
feat/distributed-map
Draft

nvasiu wants to merge 2 commits into
mainfrom
feat/distributed-map

Conversation

@nvasiu

@nvasiu nvasiu commented Aug 28, 2026 •

Copy link
Copy Markdown
Owner

By submitting this pull request, I confirm that my contribution is made under the terms of the Apache 2.0 license.

Issue Link, if available

N/A

Description

Adds the distributed map operation (ctx.distributedMap) to the Java SDK.

A map run processes a bounded dataset in parallel. A customer starts a map run from a durable function, naming a source to read items from, a processor function to invoke per batch, and concurrency, retry, and failure settings. The service reads items from the source, groups them into batches, invokes the processor for each batch, retries failures, tracks progress, routes successful results and failed items to destinations, and reports completion.

Changes:

model/

The result types a customer receives back from a map run.

  • DistributedMapSummary: what ctx.distributedMap returns, describes the run's overall outcome. throwIfError() opts into raising.
  • DistributedMapResult: returned when a result type is passed. Contains individual map run item outcomes.
  • DistributedMapResult.Item and DistributedMapResult.ItemError: represent a single item's result / error. Nested, since each belongs to exactly one result.
  • DistributedMapStatus and DistributedMapCompletionReason: the run enums.

config/

The input types a customer constructs to describe a map run.

  • DistributedMapConfig: optional settings for distributed map, nesting CompletionConfig for the item failure thresholds that mark the overall run failed.
  • DistributedMapSource: describes where map run items come from, with factories for inline, S3 and reader sources, and nesting the CSV options.
  • DistributedMapProcessor: describes the Lambda that processes items, how outcomes are reported back via batch, itemFailures and itemResults, and how failing items are retried via maxRetryAttempts and maxRetryDuration.
  • DistributedMapDestination: routes successful and failed item records to S3, nesting Success and Failure so a success destination cannot be used for failures.

DurableContext.java and context/DurableContextImpl.java

The customer-facing entry point on the durable execution context.

  • ctx.distributedMap: the method a customer calls to run a distributed map. Overloaded so the return type is DistributedMapResult when a result type is passed and DistributedMapSummary otherwise, with Class<O> and TypeToken<O> forms, blocking and async forms, and optional config.
  • Validates maxConcurrency against the model range of 1 to 10000.

dmap/

Authoring helpers for the processor and reader Lambdas, so a customer writes a plain function rather than the item or batch protocol.

  • DistributedMapHandlers: createDistributedMapItemHandler, createDistributedMapBatchHandler, createDistributedMapReader, and the durable variants createDistributedMapItemHandlerWithDurableExecution and createDistributedMapBatchHandlerWithDurableExecution.
  • The item handlers run at most concurrency handlers at a time, defaulting to one.
  • ProcessorEvent, ProcessorRecord, ItemHandlerResponse, ItemResult, ItemFailure, ReaderEvent, ReaderResponse: the wire envelopes as records, so the handler signatures name concrete types and the Lambda runtime handles serialization. Envelope validation lives in their compact constructors, which the runtime invokes when it deserializes.
  • ReaderPage: what a customer's reader function returns, items plus the next state.

operation/DistributedMapOperation.java

The executor that drives the operation against the durable execution runtime.

  • Suspends the caller's function while the run executes and resumes it with the finished outcome.
  • Every terminal state resolves, and throwIfError() opts into raising.
  • All translation between the config types and the service shapes happens here, mapping onto DistributedMapOptions and parsing DistributedMapDetails back.
  • Items, Output and the record body are opaque strings, matching the API model, so any serdes works.

operation/ChildContextOperation.java

Operation type handling.

  • Handles DISTRIBUTED_MAP in the operation sub-type switch.

exception/DistributedMapException.java

The error type a customer catches.

  • Thrown when a run or an item fails, carrying the run status, completion reason and failure count.

util/DistributedMapValidation.java

Shared validation and S3 URI parsing.

execution/CheckpointManager.java

Batching of operation updates into checkpoint requests.

  • estimateSize() now counts a distributed map START update's inline source items, so a large inline source is no longer measured as about 100 bytes when updates are packed into a checkpoint request.

Tests

sdk/src/test/.../config/

  • Config and argument validation: DistributedMapConfigTest (11), DistributedMapSourceTest (36), DistributedMapProcessorTest (17), DistributedMapDestinationTest (14).

sdk/src/test/.../model/

  • The result types and enums: DistributedMapResultTest (11), DistributedMapSummaryTest (7), DistributedMapStatusTest (4), DistributedMapCompletionReasonTest (4).

sdk/src/test/.../operation/

  • Executor behaviour and serialization round trips: DistributedMapOperationTest (8), DistributedMapOperationTranslationTest (47), which asserts the wire shape in both directions.

sdk/src/test/.../context/DurableContextDistributedMapTest.java (9)

  • The ctx.distributedMap surface and its validation.

sdk/src/test/.../dmap/DistributedMapHandlersTest.java (24)

  • Authoring helper tests: checking that they process items, report failures against the right item, reject bad inputs, and emit the exact JSON the Lambda runtime sends. The concurrency tests use a barrier, so they fail if the pool does not genuinely run items in parallel.

NOTE: This branch does not compile from a clean checkout, since the distributed map shapes haven't been released to a public AWS SDK version. I build and test against a locally patched Lambda client that carries them. The pin needs bumping to the first released version with those shapes before this can merge.

Demo/Screenshots

N/A

Checklist

  • I have filled out every section of the PR template
  • I have thoroughly tested this change

Testing

Unit Tests

Have unit tests been written for these changes?

See above.

Integration Tests

Have integration tests been written for these changes?

TODO

Examples

Has a new example been added for the change? (if applicable)

TODO

Future Tasks

  • Add distributed map to the local emulator, then add full end to end tests against it.
  • Add distributed map examples to the examples package.

Comment thread sdk/src/main/java/software/amazon/lambda/durable/dmap/DistributedMapHandlers.java Outdated
Comment thread sdk/src/main/java/software/amazon/lambda/durable/DistributedMapHandlers.java Outdated
Comment thread sdk/src/main/java/software/amazon/lambda/durable/DistributedMapHandlers.java Outdated
Comment thread sdk/src/main/java/software/amazon/lambda/durable/model/DistributedMapResult.java Outdated
Comment thread sdk/src/main/java/software/amazon/lambda/durable/DistributedMapHandlers.java Outdated
@nvasiu

nvasiu commented Sep 25, 2026

Copy link
Copy Markdown
Owner Author

Refactored to address @zhongkechen comments. The following changes were made:

Fixes

  • Terminal states resolve with the summary, throwIfError() is the opt-in
  • Removed double serialization at three sites (inline items, toItem, output)
  • maxConcurrency validates 1 to 10000
  • Run status comes from the operation, parseStatus deleted
  • Durable item handler serializes before checkpointing, so a custom resultSerDes survives a suspension
  • Output serialization moved inside each task, so one bad result fails its own item
  • Non-positive concurrency rejected, new overload defaults to one
  • getResults() preserves null outputs

Deletions

  • DistributedMapWire, folded into DistributedMapOperation
  • ProcessorRetryConfig, folded onto DistributedMapProcessor
  • The original DistributedMapException, was never thrown
  • distributedMapId from the summary and result, removed because it's not useful
  • toOutputBody, four from(...) parsers, five toMap(...) builders

Renames

  • DistributedMapError to DistributedMapException
  • reportBatchOutcome / reportFailedItems / reportItemResults to batch / itemFailures / itemResults, wire strings unchanged

Nesting, eight top-level names removed

  • DistributedMapResultItem and DistributedMapItemError into DistributedMapResult
  • DistributedMapCompletionConfig into DistributedMapConfig
  • CsvFormat and CsvDelimiter into DistributedMapSource
  • SuccessDestination, FailureDestination and DistributedMapDestinationConfig collapsed into DistributedMapDestination

Package and types

  • DistributedMapHandlers from the root package, and ReaderPage from model/, into dmap/
  • Seven envelope records added in dmap/, handler signatures name them instead of Map and Object
  • ItemFailure nests ErrorDetail, ReaderEvent.maxItems boxed, ReaderResponse generic
  • Envelope validation in compact constructors
  • DISTRIBUTED_MAP added to the ChildContextOperation switch

Tests

  • Inverted two tests that asserted the old null filtering
  • Folded DistributedMapItemErrorTest and CsvFormatTest into their owners
  • DistributedMapCompletionConfigTest to DistributedMapConfigTest, DistributedMapWireTest to DistributedMapOperationTranslationTest
  • New: per-item serialization boundary, parameterized non-positive concurrency, peak-in-flight default, barrier-based parallelism, three wire-shape tests

@nvasiu
nvasiu force-pushed the feat/distributed-map branch from 837e283 to f6b8990 Compare September 28, 2026 21:53
var index = i;
futures.add(pool.submit(() -> {
I item = toItem(inSerdes, records.get(index).body(), itemType);
outputs[index] = outSerdes.serialize(func.apply(item));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is the output used for failed items?

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No, outputs is only supposed to contain the results of succeeded items. If the serialization fails here, it will throw and the serialization error will be added to errors instead.

But there is a related issue here. In the ITEM_FAILURES mode, we don't need to return outputs but we still try to serialize each item. This was causing successful items to become failed.

Adding an update to skip this unused serialization when in ITEM_FAILURES mode.

return (event, context) -> {
S state = event.state() != null ? serdes.deserialize(event.state(), stateType) : null;

var page = func.apply(state);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

how would the reader know maxItems?

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is actually an oversight, thanks for pointing it out. The service is supposed to pass the max items per page when invoking the reader function. I will update this to accept the max items.

for (var record : records) {
bodies.add(record.body());
}
MapResult<String> batch = ctx.map(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

How many threads would be created at maximum for this map operation?

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One thread per item + one thread for ctx.map itself. Concurrency is currently unbounded in this durable helper, so a thread could be started for every item in the batch (max size 10,000).

I will add a max concurrency option to this helper, and give it a default value of 1 (matching the non-durable version of this helper).

@Override
protected void start() {
sendOperationUpdate(
OperationUpdate.builder().action(OperationAction.START).distributedMapOptions(options));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible to exceed the checkpoint size limit here?

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dmap options can include an inline list of items, with a max size of 1MB. This won't exceed the checkpoint size limit. But there is a related issue:

This SDK batches checkpoint updates together, adding updates while the batch stays under 750kb, and estimateSize() is what measures each update.

But estimateSize() doesn't look at an operation's options, because before dmap no customer-controlled data went into them. So an inline list of up to 1MB will be ignored when batching, and we can end up with batches much larger than we expect.

I will update estimateSize() to check dmap's inline items.

@nvasiu

nvasiu commented Oct 7, 2026 •

Copy link
Copy Markdown
Owner Author

New revision to address recent comments:

  • dmap/DistributedMapHandlers.java: stopped serializing item outputs when the processor reports ITEM_FAILURES.
  • dmap/DistributedMapHandlers.java: passed the page bound to the reader function, which now takes BiFunction<S, Integer, ReaderPage<I, S>>.
  • dmap/ReaderEvent.java: renamed maxItems to maxItemsPerPage, matching the key the service sends.
  • dmap/DistributedMapHandlers.java: added a concurrency parameter to createDistributedMapItemHandlerWithDurableExecution, defaulting to 1 and validated by the same check the non-durable handler uses.
  • execution/CheckpointManager.java: estimateSize() counts a distributed map START update's inline source items.

Other changes:

  • config/DistributedMapSource.java: renamed maxItemsToRead to maxItems.
  • model/DistributedMapResult.java: renamed getResults to results and getErrors to errors.
  • model/DistributedMapResult.java: dropped the null filter from errors, so it returns one entry per failed item, including items that failed without error details.
  • exception/DistributedMapException.java: replaced the runLevel and itemLevel factories with two public constructors, matching how the sibling exceptions are built.
  • dmap/ItemFailure.java: trimmed the javadoc.
  • dmap/DistributedMapHandlers.java: added an itemSerDes parameter to createDistributedMapReader and serialized each item of a reader page with it, so the processor deserializes item bodies with the serdes that encoded them. ReaderResponse.items is now List<String>.
  • Tests updated for the renames, plus new coverage for the reader page bound, the durable handler's concurrency validation, the inline source's effect on batch packing, and the reader-to-processor round trip under a custom item serdes.

@nvasiu
nvasiu force-pushed the feat/distributed-map branch from f6b8990 to 74d250b Compare October 7, 2026 18:50
Addresses review feedback on the unreleased distributed map feature.

Fixes:

- The durable item handler checkpointed item results with the default
  serdes, so a custom resultSerDes produced different JSON after a
  suspension. Results are now serialized before being checkpointed.
- Output serialization ran after the per-item futures settled, so one
  unserializable result failed the whole batch. It now runs inside each
  submitted task, so the failure is reported against its own item.
- Non-positive concurrency fell back to one thread per record. It is now
  rejected, and a new overload defaults to one.
- getResults filtered out null outputs, so the list no longer lined up
  with succeeded(). It now preserves them, matching MapResult.

Changes:

- Authoring helpers move into a dmap package, keeping feature specific
  classes off the root package.
- Handler signatures use concrete records instead of Map and Object. The
  seven envelope types are public, and envelope validation moves into
  their compact constructors.
- DistributedMapWire folds into DistributedMapOperation, and
  ProcessorRetryConfig folds onto DistributedMapProcessor.
- DistributedMapError becomes DistributedMapException, matching the other
  21 exception types.
- Types owned by one aggregate are nested into it, removing 8 top level
  names.
- Adds DISTRIBUTED_MAP to the exhaustive switch in ChildContextOperation.
@nvasiu
nvasiu force-pushed the feat/distributed-map branch from 74d250b to e299101 Compare October 7, 2026 23:15
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants