feat: allow concurrent artifact uploads - #8437
richardgetz wants to merge 22 commits into
Conversation
|
@richardgetz thank you for opening this PR! Looks like there are some failing tests, could you please check them? |
…sam-cli into rick/concurrent-uploads
Resolve conflicts with develop's deploy/package refactors (JSON output mode, language-extensions export paths): thread parallel_upload through the new DeployContext/GuidedContext/PackageContext signatures and all three nested-stack Template exports. Drop the SAM_CLI_DOCKER_TIMEOUT default change from this PR (unrelated to parallel uploads). Regenerate schema, apply black, and pass parallel_upload in develop's new deploy/package unit tests. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
|
||
| [default.deploy.parameters] | ||
| stack_name="using_config_file" | ||
|
|
There was a problem hiding this comment.
[GENERAL] The new section is inserted inside the opening fenced block that starts at line 6 and closes at line 28. The nested terminates the outer fence early, so the remainder of the original example (`capabilities="CAPABILITY_IAM"`, `region=...`, `profile=...`, and the closing fence) is left dangling and the rest of the document's fencing is inverted. Move the new heading and example after the closing on line 28. The diff also strips the trailing newline at end of file.
There was a problem hiding this comment.
Fixed in 432d67c: moved the parallel_upload section after the closing fence of the sample config, so the original example is intact. The trailing newline is preserved.
| normalize_parameters=True, | ||
| parent_stack_id=resource_id, | ||
| language_extensions_enabled=False, | ||
| parallel_upload=getattr(self, "parallel_upload", False), |
There was a problem hiding this comment.
[RESOURCE_MANAGEMENT] parallel_upload is propagated into nested-stack child Template instances (here and at the two other getattr(self, "parallel_upload", False) sites). Since CloudFormationStackResource.do_export already runs inside a parent worker thread, each nested stack creates its own ThreadPoolExecutor with up to DEFAULT_PARALLEL_UPLOAD_WORKERS (32) threads while holding a parent worker. Thread count multiplies per nesting level, so a deeply nested stack can spawn hundreds of threads and saturate the network/Docker daemon. Consider a single module-level executor shared by the whole export (or a semaphore bounding total in-flight uploads) instead of one pool per Template.
There was a problem hiding this comment.
Fixed in 432d67c: added a module-level BoundedSemaphore (_UPLOAD_SLOTS, sized DEFAULT_PARALLEL_UPLOAD_WORKERS) that every leaf artifact upload acquires, so in-flight uploads stay bounded across the whole export however deep the nesting. Nested-stack exports (CloudFormationStackResource and subclasses) only coordinate their children and do not take a slot, so a parent waiting on its children can't starve them. Covered by test_parallel_export_job_bounds_leaf_uploads_but_not_nested_stacks.
| DEFAULT_PARALLEL_UPLOAD_WORKERS = max(4, min(32, (os.cpu_count() or 1) * 2)) | ||
|
|
||
|
|
||
| class _ThreadSafeUploadCache(MutableMapping[str, str]): |
There was a problem hiding this comment.
[CONCURRENCY] _ThreadSafeUploadCache makes individual __getitem__/__setitem__ calls atomic, but the consumer in upload_local_artifacts (samcli/lib/package/utils.py:190) is a check-then-act:
if previously_uploaded and local_path in previously_uploaded:
return previously_uploaded[local_path]
...
result = zip_and_upload(...)
previously_uploaded[local_path] = resultTwo jobs referencing the same CodeUri can both miss and then both zip and upload the same artifact, which defeats the dedup the cache exists for (the PackagePerformance path is exactly the case where this matters). A per-key lock, or storing a Future/placeholder under the lock so the second caller waits on the first upload, would make the dedup hold under concurrency.
There was a problem hiding this comment.
Fixed in 432d67c: _ThreadSafeUploadCache now exposes key_lock(key), and upload_local_artifacts holds that per-path lock around the check / zip-and-upload / store sequence. Concurrent jobs sharing a CodeUri upload it once; different paths still proceed in parallel. Covered by test_upload_local_artifacts_uploads_shared_path_once_across_threads.
| if not self.login_session_active: | ||
| self.login() | ||
| self.login_session_active = True | ||
| with self._login_lock: |
There was a problem hiding this comment.
[CONCURRENCY] The login double-check is now guarded, but the rest of upload is not, and it is reachable from multiple threads with a single shared ECRUploader. self.log_streamer.stream_progress writes cursor-relative ANSI sequences (cursor up N / cursor down N, computed from a per-call ids map) to the shared self.stream; concurrent pushes will move the cursor based on each other's line offsets and produce corrupted terminal output. S3Uploader's ProgressPercentage has the same problem with its \r-based writes to stderr. Either serialize progress rendering with a lock around the stream, or force no_progressbar when parallel_upload is enabled.
There was a problem hiding this comment.
Fixed in 432d67c: PackageContext now passes no_progressbar=True to both S3Uploader and ECRUploader when parallel_upload is enabled, since their in-place redraws can't be shared safely across threads. Covered by test_parallel_upload_disables_progress_bars.
| exporter.parallel_upload = self.parallel_upload | ||
| exporter.export(resource_full_path, resource_dict, self.template_dir) | ||
|
|
||
| return _job |
There was a problem hiding this comment.
[ERROR_HANDLING] wait(futures, return_when=FIRST_EXCEPTION) suggests fail-fast, but nothing is cancelled: the enclosing with ThreadPoolExecutor(...) calls shutdown(wait=True) without cancel_futures=True, so every queued job still runs before the exception surfaces. Additionally, for future in futures: future.result() raises on the first failure in submission order, so exceptions from other failed jobs are never retrieved or logged and the user sees only one of several failures. Consider executor.shutdown(cancel_futures=True) on first error and logging the remaining exceptions before re-raising.
There was a problem hiding this comment.
Fixed in 432d67c: on the first failure the executor is shut down with cancel_futures=True (queued jobs are dropped, running ones finish), every other failure is logged, and the first failure to complete is re-raised. Covered by test_execute_jobs_in_parallel_fails_fast_and_reports_all_failures.
- Bound in-flight uploads across nested stacks with a shared semaphore; nested-stack exports coordinate children and do not take a slot. - Serialize check-then-upload per local path so concurrent jobs sharing a CodeUri upload it once (per-key lock on the thread-safe cache). - Disable S3/ECR progress bars when uploading in parallel; their in-place redraws interleave across threads. - Fail fast: cancel queued jobs on the first failure and log every other failure before re-raising the first. - Move the samconfig docs example outside the enclosing code fence. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| parallel_upload = True | ||
| else: | ||
| click.secho("\t#Speed up artifact uploads by running them in parallel") | ||
| parallel_upload = confirm(f"\t{self.start_bold}Enable parallel uploads{self.end_bold}", default=False) |
There was a problem hiding this comment.
[GENERAL] The new confirm is inserted into the middle of the guided flow (after "Disable rollback", before prompt_authorization), which shifts every subsequent answer in the guided-deploy integration tests. Those tests feed answers positionally over stdin, e.g. tests/integration/deploy/test_deploy_command.py:766:
f"{stack_name}\n\n\n\n\ny\n\n\n{autorization_question_answer}n\n{self.ecr_repo_name}\n\n\n\n".encode(),With one extra prompt ahead of the authorization question, the n and {self.ecr_repo_name} values are consumed by the wrong prompts, so the repo name lands in the authorization/ECR-deletion answer. The same shift affects lines 707, 734, 798, 824 and tests/integration/deploy/test_managed_stack_deploy.py:78. The unit tests in test_guided_context.py were updated, but the checklist marks integration tests NA — these need updating too (add one \n at the right position), which likely explains the failing tests flagged in review.
Also note the branch skips the prompt entirely when --parallel-upload is passed, unlike every neighbouring option which passes the incoming value as default= (e.g. default=self.disable_rollback). Using the existing pattern would both keep the flow consistent and keep the prompt count stable:
parallel_upload = confirm(
f"\t{self.start_bold}Enable parallel uploads{self.end_bold}", default=self.parallel_upload
)There was a problem hiding this comment.
Fixed in 95d5dbc by removing the new guided prompt rather than shifting the integration-test inputs: --guided now asks exactly the same questions as on develop (the unit tests in test_guided_context.py and the guided cases in test_command.py are back to their develop versions and pass unmodified), so scripted guided deploys and the positional stdin answers in tests/integration/deploy/ are unaffected. The --parallel-upload flag value is still carried through --guided and saved to samconfig (test_guided_prompts_keep_parallel_upload_flag_without_prompting).
| jobs[0]() | ||
| return | ||
|
|
||
| executor = ThreadPoolExecutor(max_workers=max_workers) |
There was a problem hiding this comment.
[RESOURCE_MANAGEMENT] Re-raising the nested-stack thread multiplication issue from the previous round. _UPLOAD_SLOTS correctly bounds in-flight uploads and avoids parent/child starvation, but it does not bound threads: parallel_upload is still propagated into every nested-stack child Template, and each one enters _execute_jobs_in_parallel and constructs its own ThreadPoolExecutor. Since all jobs are submitted immediately, each executor eagerly spawns up to min(len(jobs), DEFAULT_PARALLEL_UPLOAD_WORKERS) threads, most of which then just block on the semaphore. For a stack with many nested stacks the total live thread count is the sum across all templates (hundreds), not 32.
A single shared executor threaded through the Template constructor (created once at the root, reused by children) would bound both threads and uploads with one mechanism and let you drop _UPLOAD_SLOTS entirely. If you prefer to keep per-template executors, size them against a remaining-thread budget rather than the global constant.
There was a problem hiding this comment.
Fixed in 95d5dbc along the lines you suggested: the root Template creates a single ThreadPoolExecutor (DEFAULT_PARALLEL_UPLOAD_WORKERS) and passes it to every nested-stack child Template via a new upload_executor argument, so one pool bounds both threads and in-flight uploads for the whole export; the _UPLOAD_SLOTS semaphore is gone. Leaf uploads from every level are submitted to that pool. Nested-stack exports run in the calling thread instead of a pool worker, because a nested stack waits on its children and doing that from inside a shared pool could deadlock it. If a nested-stack export fails, pending uploads for that template are cancelled. Covered by test_run_export_jobs_submits_uploads_and_runs_nested_stacks_inline, test_run_export_jobs_cancels_uploads_when_nested_stack_fails, and test_export_job_passes_shared_executor_to_nested_stacks.
| if failed: | ||
| # Fail fast: drop queued jobs, let running ones finish, then report every failure. | ||
| executor.shutdown(wait=True, cancel_futures=True) | ||
| first_error = cast(BaseException, failed[0].exception()) |
There was a problem hiding this comment.
[ERROR_HANDLING] done returned by wait() is a set, so failed has arbitrary ordering. When more than one job has already failed by the time wait() returns (common: several uploads fail against the same bad bucket/credentials), failed[0] — and therefore the exception the user actually sees — is nondeterministic across runs, while the rest are only logged. The comment "report every failure" holds, but which one surfaces does not.
Filter in submission order so the propagated error is stable:
failed = [future for future in futures if future in done and future.exception() is not None]Other items from the prior review verified as addressed: the docs/sam-config-docs.md fence is now placed after the closing fence; the cache check-then-act is serialized via _ThreadSafeUploadCache.key_lock; queued jobs are now cancelled via cancel_futures=True; and the ECRUploader cursor-corruption concern is neutralized because PackageContext forces no_progressbar when parallel_upload is set, routing each push through a per-call LogStreamer on a private StringIO.
There was a problem hiding this comment.
Fixed in 95d5dbc using your suggestion: failed is now filtered in submission order ([f for f in futures if f in done and f.exception() is not None]), so the propagated error is deterministic; the others are still logged. Covered by test_wait_fail_fast_surfaces_first_failure_in_submission_order.
- Create a single ThreadPoolExecutor at the root Template and pass it to nested-stack child Templates, so one pool bounds both threads and in-flight uploads for the whole export (replaces the upload semaphore). Nested-stack exports run in the calling thread rather than a pool worker, since they wait on their children; their children's uploads go to the shared pool. - Surface the first failure in submission order so the reported error is deterministic when several uploads fail; cancel pending uploads if a nested-stack export fails. - Remove the new --guided prompt: it shifted the positional answers used by scripted guided deploys and the guided integration tests. --guided keeps the --parallel-upload flag value and saves it to samconfig. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| cache: Optional[MutableMapping[str, str]] = None | ||
| if is_experimental_enabled(ExperimentalFlag.PackagePerformance): | ||
| cache = {} | ||
| if cache is not None and self.parallel_upload: |
There was a problem hiding this comment.
[PERFORMANCE] The thread-safe dedup cache is only created when the PackagePerformance experimental flag is on:
cache: Optional[MutableMapping[str, str]] = None
if is_experimental_enabled(ExperimentalFlag.PackagePerformance):
cache = {}
if cache is not None and self.parallel_upload:
cache = _ThreadSafeUploadCache(cache)On the default path (experimental flag off) cache stays None, so upload_local_artifacts gets previously_uploaded=None, key_lock is None, and the per-path lock degrades to nullcontext(). Serial execution still deduplicated implicitly, because S3Uploader.upload does if not self.force_upload and self.file_exists(remote_path) and the second resource sharing a CodeUri hit an object that the first had already finished uploading. Under --parallel-upload all N jobs run dir_checksum + zip_folder + file_exists before any upload completes, so every one of them zips and uploads the same bytes. A template with several functions pointing at one shared code directory (a common layout) will do N× the zipping and N× the upload traffic with the flag on — the opposite of the flag's purpose.
Wrapping unconditionally when parallel uploads are on restores per-path serialization independently of the experimental flag:
if self.parallel_upload:
cache = _ThreadSafeUploadCache(cache or {})Note this also makes --parallel-upload start populating the dedup cache outside the experimental flag, so it is worth confirming that is the intended coupling.
There was a problem hiding this comment.
Good catch on the concurrent-duplicate case, fixed in 69662de, but with a narrower change than suggested. One correction on the cost: serial export already zips once per resource (zip_and_upload zips before upload_with_dedup); what serial avoids is the upload, via the file_exists check. So the parallel regression is duplicate uploads of identical content, not extra zipping.
I deliberately did not make --parallel-upload populate the dedup cache outside the PackagePerformance flag. That cache is keyed only on the local path, so if one directory is used by resources that package it differently (e.g. a function zipped with make_zip_with_lambda_permissions and a layer zipped with make_zip), the second resource would receive the first one's S3 URL. Coupling that to a performance flag seemed like the wrong trade.
Instead, in parallel mode the export now always passes a _ThreadSafeUploadCache; without the experimental flag it is lock-only (store=False: results are not remembered). Jobs sharing a local path take the same per-path lock and run one after another, so each still zips (as today) and the uploader's existence check skips re-uploading identical content, which is exactly serial behaviour. Different paths still upload concurrently. Covered by test_template_export_parallel_without_experimental_cache_uses_lock_only_cache, test_lock_only_upload_cache_does_not_remember_results and test_upload_local_artifacts_serializes_shared_path_with_lock_only_cache.
| except BaseException: | ||
| for future in futures: | ||
| future.cancel() | ||
| wait(futures) |
There was a problem hiding this comment.
[ERROR_HANDLING] _wait_fail_fast logs every non-surfaced failure, but the nested-stack failure path in _run_export_jobs does not:
except BaseException:
for future in futures:
future.cancel()
wait(futures)
raise_wait_fail_fast is never reached here, so any leaf upload future that already failed has its exception discarded without ever being read or logged — the user only sees the nested-stack error. Since a shared bad bucket or expired credentials typically fails the nested-stack template upload and the sibling artifact uploads together, this is the case where the extra detail matters most. Logging the collected upload failures before re-raising would make both paths consistent:
except BaseException:
for future in futures:
future.cancel()
wait(futures)
for future in futures:
if future.cancelled():
continue
error = future.exception()
if error is not None:
LOG.error("Parallel artifact upload also failed: %s", error, exc_info=error)
raiseSkipped as generated/derived: schema/samcli.json (regenerated from the option help text; the added parallel_upload entry matches the new flag).
Prior-round findings verified as resolved in this revision: the docs/sam-config-docs.md fence now closes before the new section; nested stacks share a single ThreadPoolExecutor instead of creating one each; upload_local_artifacts holds a per-path lock around check/zip/store; PackageContext forces no_progressbar for both uploaders, and ECRUploader then streams into a per-call StringIO rather than the shared stream; queued futures are cancelled on first failure and the surfaced error is stable in submission order; the guided-deploy prompt was dropped entirely, so stdin answer ordering in the integration tests is unchanged.
There was a problem hiding this comment.
Fixed in 69662de: factored the logging into Template._log_other_failures and call it on the nested-stack failure path too (after cancelling and waiting), so leaf upload failures are logged there as well and both paths behave the same. Covered by test_run_export_jobs_logs_upload_failures_when_nested_stack_fails.
…lures - In parallel mode always pass a thread-safe cache so jobs sharing a local path take the same per-path lock. Without the PackagePerformance flag the cache is lock-only (results are not stored), so behaviour matches a serial export: every job zips, and the uploader's remote existence check skips re-uploading identical content. This avoids enabling the path-keyed experimental cache implicitly. - Log leaf upload failures on the nested-stack failure path as well, so no upload error is dropped. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| """ | ||
| futures = [executor.submit(job) for is_nested_stack, job in jobs if not is_nested_stack] | ||
| try: | ||
| for is_nested_stack, job in jobs: |
There was a problem hiding this comment.
[ERROR_HANDLING] The inline nested-stack loop never checks whether an already-submitted upload has failed, so fail-fast only applies once the loop finishes:
futures = [executor.submit(job) for is_nested_stack, job in jobs if not is_nested_stack]
try:
for is_nested_stack, job in jobs:
if is_nested_stack:
job()
except BaseException:
...
Template._wait_fail_fast(futures)If a leaf upload fails for a reason that dooms the whole export (expired credentials, NoSuchBucket, a bad --image-repository), the loop keeps running nested-stack exports to completion — each of which submits and then waits on its own children's uploads, all of which will also fail against the same bucket/registry. Only after the last nested stack returns does _wait_fail_fast surface the original error. Serial export aborts at the first failing resource, so this is a behavior regression: many doomed uploads are attempted and the user waits for the full tree before seeing the error.
Checking for a completed-with-exception future between nested-stack jobs would restore the fail-fast contract, e.g.:
for is_nested_stack, job in jobs:
if not is_nested_stack:
continue
if any(f.done() and not f.cancelled() and f.exception() is not None for f in futures):
break # fall through to _wait_fail_fast, which cancels and re-raises
job()There was a problem hiding this comment.
Fixed in 84e9cae as part of the nested-stack change (see the PERFORMANCE thread): nested-stack exports are no longer run in a loop after the uploads are submitted. All of a template's uploads and nested-stack exports are started together and a single _wait_fail_fast covers both, so the first failure — upload or nested stack — surfaces immediately; pending uploads are cancelled and every other failure is logged. There is no longer a point where more nested stacks are started after a failure is known. Covered by test_run_export_jobs_reports_nested_stack_and_upload_failures and test_run_export_jobs_cancels_uploads_when_nested_stack_fails.
| try: | ||
| for is_nested_stack, job in jobs: | ||
| if is_nested_stack: | ||
| job() |
There was a problem hiding this comment.
[PERFORMANCE] Nested-stack exports run inline in the calling thread, one at a time, and each child blocks on all of its own uploads before the next child starts: the child's export() enters _run_export_jobs, submits its leaf jobs to the shared pool, and ends with _wait_fail_fast(futures); CloudFormationStackResource.do_export then uploads the child template; only then does the parent loop advance to the next nested stack.
The consequence is that child N+1's uploads are not even submitted until child N's have all completed, and every child's uploads are queued behind the parent's own leaf uploads. For a template whose artifacts live mostly inside nested stacks — the "larger stacks" case in #7056 — the concurrency window is one nested stack's worth of leaf artifacts at a time, so --parallel-upload buys much less than the flag suggests.
Two ways to address it: collect leaf jobs across the whole tree and submit them from a single wait point at the root, or run nested-stack coordination on a separate small executor (coordination threads only wait, they don't upload) so several children can submit into the shared upload pool at once. If the current shape is intentional for this iteration, the flag's help text should say that parallelism applies per nested stack rather than across the whole tree.
There was a problem hiding this comment.
Agreed, fixed in 84e9cae with the second approach you suggested. Nested-stack coordination now runs on a small separate pool (one thread per nested stack at that level) rather than inline, so sibling nested stacks run concurrently and all submit their leaf uploads into the single shared upload pool; the coordination threads only wait and never occupy an upload worker, so the shared pool cannot deadlock. Total threads are the upload pool plus one per nested stack in the template, which is fixed by template structure rather than multiplied per level. Covered by test_run_export_jobs_runs_sibling_nested_stacks_concurrently (two nested stacks that each wait on the other only complete if they run at the same time) and test_run_export_jobs_runs_nested_stacks_outside_the_upload_pool.
|
|
||
| # NOTE: sriram-mv, A cyclic dependency on `Template` needs to be broken. | ||
|
|
||
| DEFAULT_PARALLEL_UPLOAD_WORKERS = max(4, min(32, (os.cpu_count() or 1) * 2)) |
There was a problem hiding this comment.
[RESOURCE_MANAGEMENT] DEFAULT_PARALLEL_UPLOAD_WORKERS = max(4, min(32, (os.cpu_count() or 1) * 2)) is the only control over concurrency, and it is not user-configurable. Every leaf ZIP artifact goes through zip_and_upload → zip_folder, which writes a full zip into a fresh tempfile.mkdtemp() before uploading and only removes it after the upload finishes. Serial export holds at most one such zip on disk; with this pool up to 32 exist simultaneously, so peak temp-directory usage scales with worker count times artifact size — on a 16-core CI runner with several large Lambda packages that is gigabytes of /tmp that was previously never needed. make_zip_with_permissions also buffers each member file fully in memory (file_bytes = data.read()), multiplying peak RSS by the same factor.
Since --parallel-upload is a flag rather than a value, users who hit this have no way to dial it back. Consider accepting an optional worker count (e.g. --parallel-upload as an int option, or an accompanying env var) and/or choosing a lower default cap, since artifact upload is I/O-bound and does not need cpu_count * 2 workers to saturate bandwidth.
There was a problem hiding this comment.
Fixed in 84e9cae: the default is now a fixed 8 concurrent uploads rather than cpu_count * 2 up to 32, since uploads are I/O bound and each in-flight zip holds a temp file, and it can be changed with the SAM_CLI_PARALLEL_UPLOAD_WORKERS environment variable (invalid or < 1 values fall back to the default with a warning). I kept --parallel-upload a boolean flag and used an env var, in the same style as SAM_CLI_POLL_DELAY, so samconfig parallel_upload = true keeps working; the flag help text, samconfig docs and schema describe the override. Covered by test_parallel_upload_workers_from_environment. For reference, 8 workers is what we measured against: on our ~50-image-function stacks it cut the ECR push phase from ~238s to ~21s per stack.
- Run nested-stack coordination on a small per-level pool (one thread per nested stack) instead of inline, so sibling nested stacks submit their uploads to the shared upload pool at the same time. Coordination threads only wait, so they stay out of the upload pool. One fail-fast wait now covers uploads and nested stacks, so a failure is no longer hidden behind the remaining nested-stack exports. - Default to 8 concurrent uploads (each in-flight zip upload holds a temp zip on disk) and allow overriding it with SAM_CLI_PARALLEL_UPLOAD_WORKERS; document it in the flag help, samconfig docs and schema. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| if not nested_jobs: | ||
| Template._wait_fail_fast(upload_futures) | ||
| return | ||
| with ThreadPoolExecutor(max_workers=len(nested_jobs)) as coordinator: |
There was a problem hiding this comment.
[RESOURCE_MANAGEMENT] The nested-stack coordinator pool is sized by the number of nested stacks at that level, with no cap, and a new pool is created at every level:
with ThreadPoolExecutor(max_workers=len(nested_jobs)) as coordinator:
nested_futures = [coordinator.submit(job) for job in nested_jobs]
Template._wait_fail_fast(upload_futures + nested_futures)Because all nested jobs are submitted immediately, ThreadPoolExecutor materializes one thread per job, and a parent coordinator thread stays alive for as long as its subtree is exporting. The number of live coordinator threads therefore equals the total number of nested stacks in the whole hierarchy, not per level. A template hierarchy with a few hundred AWS::CloudFormation::Stack / AWS::Serverless::Application resources produces a few hundred simultaneous threads, and SAM_CLI_PARALLEL_UPLOAD_WORKERS has no effect on that number.
This also means SAM_CLI_PARALLEL_UPLOAD_WORKERS is not an actual bound on concurrent uploads. CloudFormationStackResource.do_export uploads the rendered child template itself:
with mktempfile() as temporary_file:
...
url = self.uploader.upload(temporary_file.name, remote_path)That call runs on the coordinator thread, not in the shared upload pool, so concurrent S3 uploads can exceed the configured worker count by the number of in-flight nested stacks — each also holding a temp file on disk. Consider a single shared, bounded coordinator pool created alongside the upload pool (passed down with upload_executor) rather than one per level, and routing the nested template upload through the bounded upload pool.
There was a problem hiding this comment.
Partly addressed in 1418fe5.
Nested template uploads — fixed. CloudFormationStackResource.do_export now submits the rendered child-template upload to the shared upload pool when one is present (the coordination thread waits on it), so SAM_CLI_PARALLEL_UPLOAD_WORKERS bounds every S3 upload, including nested templates. Covered by test_nested_stack_template_upload_goes_through_shared_upload_pool.
Coordinator pool — keeping the current shape, deliberately. A single bounded coordinator pool shared across levels would deadlock: a nested-stack job holds its coordinator slot while it waits for its own nested children, and those children need slots in the same pool. With depth > 1 and more in-flight ancestors than slots, every slot is held by a parent waiting on a child that can never start. That is the same failure mode we avoided by keeping coordination out of the upload pool. The current coordinator threads only block on futures (no uploads, no temp files now that the template upload is pooled), and their count equals the number of nested stacks in the template, which is fixed by the template rather than by worker settings. I think that is the right trade versus a bounded-but-deadlock-prone pool, but happy to revisit if maintainers prefer, e.g., capping nesting concurrency by falling back to inline export past a threshold.
| if previously_uploaded and local_path in previously_uploaded: | ||
| result = previously_uploaded[local_path] | ||
| LOG.debug("Skipping upload of %s since is already uploaded to %s", local_path, result) | ||
| return cast(str, result) |
There was a problem hiding this comment.
[PERFORMANCE] The per-path lock is acquired inside a job that is already occupying an upload-pool worker, and it is held across the full zip-and-upload:
key_lock = getattr(previously_uploaded, "key_lock", None)
with key_lock(local_path) if key_lock else contextlib.nullcontext():
...
result = zip_and_upload(...)With store=False (the default path, since the experimental cache is off) nothing is ever recorded, so every waiter re-runs zip_and_upload after acquiring the lock. Templates where many functions share one CodeUri are common in SAM, and in that shape all N jobs contend on a single lock: up to workers threads sit blocked while holding pool slots, and unrelated artifacts (layers, images with distinct paths) that are still queued cannot start until those slots free up. The net effect for a mixed template is slower than the serial path, which is the opposite of the flag's purpose.
Since the mapping lives for the duration of one export and is not the cross-run experimental cache, recording the result in parallel mode regardless of store would let the second and later waiters return immediately instead of re-zipping and re-uploading, and would release the pool slots promptly. If keeping store=False is required, the lock should at least not be held by a pool worker across the upload (e.g. skip the lock and rely on the uploader's remote dedup) so distinct-path work is not starved.
There was a problem hiding this comment.
Agreed, fixed in 1418fe5. In parallel mode without the experimental cache, uploads are now memoized for the duration of the export, so the first job for a shared CodeUri zips and uploads it and the waiters return the result immediately instead of re-zipping while holding workers.
The memo is keyed by local path plus packaging: the zip method (Lambda-permission zip vs plain zip) and extension, via _ThreadSafeUploadCache(key_by_packaging=True).cache_key(...), so a directory used by a Lambda resource and by a non-Lambda resource that zips it plainly (e.g. an Elastic Beanstalk application version) stays two separate artifacts. The experimental PackagePerformance cache keeps its existing path keys. Covered by test_upload_local_artifacts_memoizes_by_packaging_across_threads and test_packaging_keyed_upload_cache_separates_packaging.
A correction to my reply in the earlier PERFORMANCE thread: I gave "a function and a layer" as the example of conflicting packaging, but layers are in LAMBDA_LOCAL_RESOURCES and are zipped the same way as functions, so that example was wrong. The realistic case is Lambda vs non-Lambda resources sharing a directory, which is rarer, but the packaging-aware key covers it either way.
- In parallel mode without the experimental cache, memoize uploads for the duration of the export, keyed by local path plus packaging (Lambda zip vs plain zip, and extension), so jobs sharing a CodeUri zip and upload once and waiters return immediately instead of re-zipping while holding pool workers. Differently packaged uses of one directory stay separate. - Upload rendered nested-stack templates through the shared upload pool so the configured worker count bounds every S3 upload. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| if not nested_jobs: | ||
| Template._wait_fail_fast(upload_futures) | ||
| return | ||
| with ThreadPoolExecutor(max_workers=len(nested_jobs)) as coordinator: |
There was a problem hiding this comment.
[ERROR_HANDLING] Fail-fast does not reach into nested-stack subtrees, so one failure still lets the whole nested tree upload to completion before the error surfaces.
with ThreadPoolExecutor(max_workers=len(nested_jobs)) as coordinator:
nested_futures = [coordinator.submit(job) for job in nested_jobs]
Template._wait_fail_fast(upload_futures + nested_futures)Because the coordinator is sized to len(nested_jobs), every nested job starts immediately and none is ever queued, so the future.cancel() loop in _wait_fail_fast can never succeed on a nested future. On the first failure _wait_fail_fast cancels only this level's still-queued upload_futures and raises; coordinator.__exit__ then runs shutdown(wait=True) and blocks until every in-flight nested export finishes. Those in-flight exports keep calling _run_export_jobs on the shared pool and submit brand-new leaf uploads (plus the rendered child-template upload in do_export) after the cancellation, and nothing cancels those.
Net effect on a multi-level template: a failure in a root-level upload still pays for the full nested subtree's uploads before the user sees the error — the opposite of the documented "on the first failure cancel the rest". The same applies recursively at each level.
The cancel-based approach can't fix this on its own, since the work is started eagerly and submitted lazily. A shared abort flag would: pass a threading.Event down alongside upload_executor, set it in _wait_fail_fast before raising, and check it at the top of _job in _build_export_job and before the submit loop in _run_export_jobs so descendants stop enqueuing new uploads once any job has failed.
This is a distinct consequence of the eager coordinator pool you chose to keep in the earlier RESOURCE_MANAGEMENT thread — flagging it separately because the correctness impact (fail-fast not holding) is independent of the thread-count concern that was dismissed.
Everything else I checked in the concurrency path holds up: no upload-pool worker ever blocks on the upload pool (nested coordination always runs off-pool, and the per-key lock in upload_local_artifacts is only ever held by a thread that already owns a worker, so it always makes progress); _wait_fail_fast orders failures deterministically by submission order and logs the non-surfaced ones on both the upload and nested-stack paths; _ThreadSafeUploadCache keys by packaging so a Lambda zip and a plain zip of the same directory can't be conflated; ECRUploader's login is double-checked under a lock and auth_config is published before login_session_active; and forcing no_progressbar=True in PackageContext does eliminate the shared-cursor corruption, since the no_progressbar branch builds a per-call LogStreamer over a private StringIO and LogStreamer writes only to its own stream.
There was a problem hiding this comment.
Agreed, fixed in f32ed53 with the shared abort flag you described. The root Template creates a threading.Event alongside the shared upload executor and both are passed down to every nested-stack child (upload_abort). _wait_fail_fast sets it on the first failure; _build_export_job jobs, the submit step in _run_export_jobs, and the rendered child-template upload in CloudFormationStackResource.do_export all check it first and raise an internal _UploadAborted instead of doing work. In-flight subtrees therefore stop enqueuing new uploads as soon as any job anywhere in the export fails, instead of uploading to completion.
Two details so the user-facing behaviour stays clean: skipped work raises _UploadAborted rather than returning, so a nested stack whose uploads were skipped cannot "succeed" with missing artifacts; and _UploadAborted is never logged as an additional failure or surfaced — the error raised is the first real failure in submission order. Covered by test_failure_stops_nested_subtrees_from_starting_new_uploads, test_wait_fail_fast_surfaces_real_failure_over_aborted_work and test_export_job_is_skipped_once_export_is_aborted.
Thanks for the thorough verification of the rest of the concurrency path.
…gnal Pass a threading.Event down with the shared upload executor. The first failure anywhere in the export sets it; export jobs, nested-stack job submission and nested template uploads that have not started yet raise an internal _UploadAborted instead of doing work, so a failure no longer lets the rest of the nested tree upload to completion. Skipped work is not logged as an additional failure, and the surfaced error is always the first real failure in submission order. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Each additional failure is logged as a single line; its traceback is only emitted at debug level. A missing bucket previously printed a full traceback for every concurrent upload that failed alongside the surfaced error. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
| failed = [future for future in futures if not future.cancelled() and future.exception() is not None] | ||
| real = [future for future in failed if not isinstance(future.exception(), _UploadAborted)] | ||
| surfaced = (real or failed)[0] | ||
| Template._log_other_failures(futures, surfaced=surfaced) | ||
| raise cast(BaseException, surfaced.exception()) |
There was a problem hiding this comment.
[ERROR_HANDLING] A nested stack that was skipped after an abort counts as a real failure, so its error can hide the actual cause
CloudFormationStackResource and ServerlessApplicationResource inherit ResourceZip.export, which wraps every exception from do_export in ExportFailedError(..., ex=ex). An _UploadAborted raised inside a nested stack's do_export therefore reaches the parent's _wait_fail_fast as an ExportFailedError. This happens in two ways: through the new abort check before the child template upload, or when the child's export() re-raises _UploadAborted because it only saw aborts.
The isinstance(future.exception(), _UploadAborted) filter does not match the wrapped error, so the skipped stack counts as "real". Take sibling nested stacks A and B, where an upload inside B fails and A is then aborted. Because A comes first in submission order, (real or failed)[0] picks A. The user then sees "Unable to upload artifact … of A resource." followed by an empty message. The actual error from B only shows up as "Parallel artifact upload also failed: …". _log_other_failures has the same gap, so it also logs these skips as errors.
Recognise wrapped aborts in both places, for example with a helper:
def _is_upload_abort(error: Optional[BaseException]) -> bool:
return isinstance(error, _UploadAborted) or isinstance(getattr(error, "ex", None), _UploadAborted)Then use _is_upload_abort(future.exception()) in the real filter and in _log_other_failures.
There was a problem hiding this comment.
Good catch, confirmed and fixed in 2860bb4. ResourceZip.export does wrap do_export failures in ExportFailedError(ex=...), so a skipped nested stack reached the parent wrapped (and once per nesting level for deeper subtrees), slipping past the isinstance check.
Added _is_upload_abort(error), which follows the wrapping chain (ex, then __cause__, with a cycle guard) rather than checking one level, and used it both in the real filter in _wait_fail_fast and in _log_other_failures. In your A/B example the surfaced error is now B's real failure and A's skip is neither surfaced nor logged. Covered by test_wait_fail_fast_ignores_aborts_wrapped_by_nested_stack_exporters (A double-wrapped and submitted first) and test_is_upload_abort_follows_wrapping.
ResourceZip.export wraps do_export failures in ExportFailedError(ex=...), so a nested stack skipped after an abort reached its parent wrapped (once per nesting level) and was treated as a real failure: it could be surfaced instead of the actual error and was logged as an extra failure. Detect aborts through the wrapping chain when choosing the surfaced error and when logging the others. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
@seshubaws I'm picking this up from Rick (I'm a collaborator on his fork). Sorry for the long gap on the failing tests. Here's where it stands:
On our image-heavy stacks (~50 image functions each), the flag cut the ECR push phase of a heavy deploy from ~238s to ~98s per stack (serial vs --parallel-upload on the same SAM base). The workflow runs for the latest commit (2860bb4) are waiting on maintainer approval. Could you approve them and take a look when you have a chance? Thanks! |
Which issue(s) does this change fix?
Why is this change necessary?
stacks. Allowing users to opt into parallel uploads lets them shorten the pre-deploy packaging phase when
bandwidth and CPU permit.
How does it address the issue?
S3/ECR uploads.
cache and guarding the ECR uploader’s login state.
What side effects does this change have?
remains the same (serial uploads) unless the flag/config setting is used.
Mandatory Checklist
PRs will only be reviewed after checklist is complete
make prpassesmake update-reproducible-reqsif dependencies were changedBy submitting this pull request, I confirm that my contribution is made under the terms of the Apache 2.0 license.