Repository navigation
Conversation
cfa93bd to
cbacaaf
Compare
|
I accidentally broke the last PR. So here's a new clean PR to add the distributed map operation with the restructure you suggested in your comment here (along with fixes to other issues you pointed out). But I made the these modifications to your proposed restructure:
I also used Thoughts on these changes? I made this graph of the new operation structure to more easily see how the components relate:
Besides the restructure, I also made the following fixes based on your comments:
|
ee54c16 to
50f5eba
Compare
Add ctx.distributed_map with inline, S3, and reader sources, the config, processor, completion, and destination types, the result types, and function-authoring helpers for item and batch handlers. Serialization lives in lambda_service.py with one dataclass per API shape. Translation happens only in the executor, so config.py has no runtime import from lambda_service. Tests are split by layer. Items, Output, and the record body are opaque strings, matching the API model, so any serdes works rather than JSON only. Every terminal state resolves with the summary, and throw_if_error() opts into raising. Missing details or a missing completion reason raise ExecutionError. The processor factories are batch, item_failures, and item_results. Source and destination factories are InlineSource, S3Source, ReaderSource, and S3Destination. Retry is max_retry_attempts and max_retry_duration on the processor. Passing DistributedMapResultConfig selects the DistributedMapResult return type statically. Only the result types and their enums are exported from the package root.
50f5eba to
1ae73e8
Compare
yaythomas
left a comment
There was a problem hiding this comment.
Great rewrite, thank you very much! Status contract, string-vs-JSON on the wire, factory names, package-root exports, the return-type union. The envelope dataclasses in dmap/handlers.py and the DistributedMapResultConfig overload are especially nice. Thank you for the second pass.
Headlines:
- Please bring back a
distributed_map_idaccessor, derived properly this time. - The module-level
_build_*helpers inoperation/dmap.pywant to befrom_*factories on the wire dataclasses (pure shape translation) or executor methods (anything that serializes). That also resolves the remainingAnyparameters. list[T]rather thantuple[T, ...]for the collection fields.
Details inline.
| unprocessed_count: int | ||
| distributed_map_run_arn: str | None = None | ||
| completion_details: str | None = None | ||
| total_count: int | None = None |
There was a problem hiding this comment.
Please could you add back a distributed_map_id derived from the ARN? The last round flagged the buggy arn.rsplit(":", 1)[-1]; the right fix is to derive it correctly rather than drop the accessor, since the run id is the handle customers need when locating a run's output objects.
Rather than a bare split, mirror DurableExecutionArn.from_arn in types.py: a frozen dataclass with an anchored regex, returning None when the ARN doesn't match. Why? A split returns something for any string, including the wrong segment for an unexpected shape; the regex validates the whole grammar and gives the parsing one home and one test.
_DISTRIBUTED_MAP_RUN_ARN_PATTERN = re.compile(
r"^(arn:[^:]*:lambda:[^:]*:[^:]*:function:[^:/]+:[^:/]+"
r"/durable-execution/[^/]+/[^/]+)/distributed-map-run/([a-z0-9]+)$"
)
@dataclass(frozen=True)
class DistributedMapRunArn:
durable_execution_arn: str
run_id: str
@classmethod
def from_arn(cls, arn: str) -> DistributedMapRunArn | None:
match = _DISTRIBUTED_MAP_RUN_ARN_PATTERN.match(arn)
if not match:
return None
return cls(durable_execution_arn=match.group(1), run_id=match.group(2))Then distributed_map_id is a property that calls from_arn and returns parsed.run_id if parsed else None.
There was a problem hiding this comment.
Will follow this approach, but I'll reuse the ARN regex in types.py to validate, so we don't have the same regex in 2 places.
| } | ||
|
|
||
|
|
||
| def _build_inline_items( |
There was a problem hiding this comment.
The translation helpers are module-level functions and they mix two kinds of work. Splitting them along that line would match the rest of the codebase:
-
Pure shape translation (customer config in, wire dataclass out, no serdes, no execution context):
_build_s3_source_config,_transform_for,_build_processor_config,_build_completion_config,_destination_entry_fields,_build_destination_config. Please make thesefrom_*classmethods on the targetlambda_servicedataclass:DistributedMapS3SourceConfig.from_source(s3: S3SourceConfig),DistributedMapProcessorConfig.from_processor(processor),DistributedMapOnSuccessConfig.from_destination(d: SuccessDestination), and so on. Why? CONTRIBUTING: "Encapsulate conversion logic in afrom_xfactory andto_xmethod on a class";ErrorObject.from_exceptionis the existing example of a wire class building itself from a domain object.lambda_serviceimportingconfigis the right direction (config is the lowest layer), so no cycle._UNLIMITED_RETRY_WIREand_RESPONSE_TYPE_FOR_MODEmove next to the processor class. The kwargs-dict splat in_destination_entry_fieldsdisappears once each subclass builds itself, and each factory takes its concrete config type, which retires theAnyparameters (CONTRIBUTING §Typing).DistributedMapCompletionConfigexists in both modules; alias one at the import site, asconcurrency/models.pydoes withBatchResult as BatchResultProtocol. -
Serialization (
_build_inline_items, the readerinitial_statebranch of_build_source_config): these callserialize()with the operation id and ARN, which only the executor has, andlambda_servicecan't importserdeswithout a cycle, so they stay in the executor. Please make them methods though: today each takesoperation_idanddurable_execution_arnas parameters the caller reads offself, and_build_distributed_map_optionstakes six parameters, all six of which areself.*, from one call site. A smallself._serialize(serdes, value)reads the two context values once. Then the wire class takes strings:DistributedMapSourceConfig.create_inline(items: list[str], max_items=),.create_s3(s3, max_items=),.create_reader(function_name, initial_state: str | None, max_items=), in the style ofOperationUpdate.create_step_start/create_invoke_start.InvokeOperationExecutor.check_result_statusis the existing example: serialize withself.*context, hand the string to the wire constructor. (serialize()already defaultsserdes=Noneto the JSON serdes, so theor DEFAULT_JSON_SERDEScan go too.) -
Result reconstruction (
_resolve_summary,_distributed_map_status_from_operation): considerDistributedMapSummary.from_operation(operation)indmap/models.py, withCheckpointedResult.create_from_operationas the precedent.DistributedMapResultadditionally needs the deserialized items, which stay an executor concern, soDistributedMapResult.from_operation(operation, items).
After this operation/dmap.py is the executor class and nothing else, which is what operation/invoke.py looks like. Mostly moving code rather than writing it.
There was a problem hiding this comment.
Thanks for the suggestions, preparing an update now with a couple changes:
serialize() already defaults serdes=None to the JSON serdes, so the or DEFAULT_JSON_SERDES can go too.
serialize() currently defaults to EXTENDED_TYPES_SERDES, so removing DEFAULT_JSON_SERDES changed the result and failed the tests. I'll need to keep that.
DistributedMapResult additionally needs the deserialized items, which stay an executor concern, so DistributedMapResult.from_operation(operation, items).
DistributedMapResult subclasses DistributedMapSummary which already has from_operation(operation). So if I try to override it with an extra required parameter (from_operation(operation, items)), mypy will reject it.
Instead I'll give DistributedMapResult a new from_operation_and_items(operation, items). And I will overload DistributedMapResult.from_operation() to return an error saying its not valid on a result, and to use from_operation_and_items() instead.
| class DistributedMapInlineSourceConfig: | ||
| """Represent the items an inline map run source reads from.""" | ||
|
|
||
| items: tuple[str, ...] = () |
There was a problem hiding this comment.
list[MyType] is probably more natural here than tuple[MyType, ...], and throughout the codebase it's mostly list[T] for this sort of thing (BatchResult.all, CheckpointUpdatedExecutionState.operations, ErrorObject.stack_trace). It's true tuple gives an immutability that could arguably belong on a frozen dataclass, but Sequence[T] / list[T] is the more common shape for a single-item-type collection. Same for the other tuple fields in config.py and dmap/handlers.py; the tuple(...) / list(self.x) conversions go with them. Pls remember to keep the copy in InlineSource.of (list(items)) so a caller mutating their list afterwards doesn't change what gets checkpointed.
|
|
||
| PASS_THROUGH_SERDES: SerDes[Any] = PassThroughSerDes() | ||
|
|
||
| _MAX_CONCURRENCY_LIMIT = 10000 |
There was a problem hiding this comment.
_MAX_CONCURRENCY_LIMIT = 10000 lives in context.py, next to nothing else about distributed map. Consider a ClassVar on DistributedMapProcessor beside UNLIMITED, or next to the wire dataclass that carries the value, so the bound sits with the thing it bounds. Trivial.
| # Replay with the run completed. | ||
| _map_run_id, replay_event = _replay_event( | ||
| { | ||
| "Status": "SUCCEEDED", |
There was a problem hiding this comment.
The DistributedMapDetails fixtures carry a "Status": ... key, but DistributedMapDetails.from_dict doesn't read one (status comes from the operation, which the fixture also sets). Harmless since the key is ignored, but it suggests a shape the code doesn't have; dropping it would make the fixtures match from_dict.
| raise ValidationError(msg) | ||
|
|
||
|
|
||
| def distributed_map_item_handler( |
There was a problem hiding this comment.
The concurrency default moved from len(records) to 1, and the docstring says so. I think that's the right default: the old one required customer code to be thread-safe out of the box, which is a surprising precondition. Just recording the deliberate behaviour change, no action needed.
| summary = self._resolve_summary(operation) | ||
| return CheckResult.create_completed(summary) | ||
|
|
||
| # Operation-level terminal failure. Every terminal state resolves with the |
There was a problem hiding this comment.
This branch is exactly right: every terminal status resolves with the summary, and ExecutionError is raised only when the backend omitted the details block. With test_terminal_failure_with_details_resolves_with_summary and test_operation_level_terminal_failure_without_details_raises the contract is well pinned. Nice!
| ) | ||
| return executor.process() | ||
|
|
||
| @overload |
There was a problem hiding this comment.
Cleanly done. config: DistributedMapResultConfig selects DistributedMapResult, everything else the summary, no isinstance at the call site. Nice!
| several threads. | ||
| """ | ||
| if func is None: | ||
| return functools.partial( |
There was a problem hiding this comment.
if func is None: return functools.partial(...) mirrors durable_execution, so bare, with-options and direct-call forms all work, and the envelope dataclasses mean no dict[str, Any] plumbing anywhere. Nice!
|
@yaythomas SummaryAdding back
Moving translation helpers:
Using list instead of tuple:
Moved _MAX_CONCURRENCY_LIMIT:
Removed status from e2e fixtures:
Other changes:
|
- The wire types translate their config-layer counterparts through from_* and create_* classmethods, so operation/dmap.py holds the executor alone and its serialization reads the operation id and execution ARN off self. - dmap/models.py rebuilds the summary and the result from a terminal operation, and parses the run ARN to expose the run id. - The source and destination factories are classmethods on the type they return, and the resolved types that are not a config= argument drop the Config suffix. - Collection fields are list rather than tuple, matching the rest of the SDK. - The item handler decorators take a typed response_mode in place of a report string. - The concurrency bound and the wire sentinels sit beside the types they apply to, and each shared type is imported from the module that defines it. - The executor fetches its checkpoint through the base class helper, which raises NonDeterministicExecutionError on an identity mismatch. - Each test file covers the module it is named for.

Summary
Adds the distributed map operation (
ctx.distributed_map) to the Python 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
dmap/models.pyThe result types a customer receives back from a map run.
DistributedMapSummary: whatctx.distributed_mapreturns, describes the run's overall outcome.DistributedMapResult: returned when aDistributedMapResultConfigis passed. Contains individual map run item outcomes.DistributedMapResultItemandDistributedMapItemError: represent a single item's result / error.dmap/__init__.pyEmpty package marker for the
dmappackage.config.pyThe input types a customer constructs to describe a map run, and the distributed map enums.
DistributedMapConfig: optional settings for distributed map.DistributedMapResultConfig: subclass that additionally collects item results inline, selecting the return type statically.InlineSource,S3Source,ReaderSource: describe where map run items come from.DistributedMapProcessor: describes the Lambda that processes items, how outcomes are reported back, and how failing items are retried viamax_retry_attemptsandmax_retry_duration.DistributedMapCompletionConfig: defines item failure thresholds for marking the overall map run failed.S3Destination,DistributedMapOnSuccessConfig,DistributedMapOnFailureConfig,DistributedMapDestinationConfig: for routing successful and failed item records to S3.lambda_service, soconfig.pyhas no runtime import from it.context.pyThe customer-facing entry point on the durable execution context.
ctx.distributed_map: the method a customer calls to run a distributed map. Overloaded so passing aDistributedMapResultConfigtypes the return asDistributedMapResultand anything else asDistributedMapSummary.max_concurrencyagainst the documented ceiling of 10000.dmap/handlers.pyAuthoring decorators for the processor Lambda, so a customer writes a plain function rather than the item or batch protocol.
distributed_map_item_handler,distributed_map_batch_handler,distributed_map_reader, and the durable variantsdurable_distributed_map_item_handleranddurable_distributed_map_batch_handler.durable_execution, so they work bare, with options, or called directly with a function.operation/dmap.pyThe executor that drives the operation against the durable execution runtime.
throw_if_error()opts into raising.lambda_service.pySerialization for carrying the operation and its results to and from the backend service.
Wiresuffix.Items,Output, and the record body are opaque strings, matching the API model, so any serdes works rather than JSON only.state.pyDurable execution state handling, so a run's outcome persists across suspend and resume.
exceptions.pyThe error type a customer catches.
DistributedMapError: raised when a run or an item fails.plugin.pyOperation type registration.
DISTRIBUTED_MAPtoOperationType.__init__.pyThe package's public API surface.
durable_executionanddurable_stepare exported.config.Tests
Tests are split by layer so each sits alongside its module.
tests/config_test.pytests/lambda_service_test.pytests/dmap/models_test.pytests/operation/dmap_test.pytests/context_test.pyctx.distributed_mapsurface.tests/dmap/handlers_test.pytests/e2e/dmap_int_test.pytests/e2e/dmap_helpers_int_test.pyFuture Tasks
By submitting this pull request, I confirm that you can use, modify, copy, and redistribute this contribution, under the terms of your choice.