Skip to content

Cancelling an async multipart copy leaves its multipart upload behind #1046

Description

@laughingman7743

Problem

When the task running AioS3FileSystem._copy_object_with_multipart_upload() is cancelled, the multipart upload is neither completed nor aborted, so the incomplete upload stays in the destination bucket (and its parts accrue storage) until a lifecycle rule removes it. This happens, for example, with asyncio.wait_for() timing out, task.cancel(), or a cancelled asyncio.gather().

  • The cleanup in pyathena/filesystem/s3_async.py (_copy_object_with_multipart_upload()) is except Exception:. asyncio.CancelledError is a BaseException, so it skips the wait for the running parts and the abort.
  • The part copies run with asyncio.to_thread(). Cancelling the awaiting tasks does not stop the threads, so parts that are already running finish after the cancellation and are stored in the upload.
  • The sync path handles interrupts: since Leave the existing object unchanged when put_file() fails #1017 (49edeea), S3FileSystem._finish_multipart_upload() catches BaseException, waits for the parts that could not be cancelled, and aborts the upload.

#973 (PR #1036) added the abort for ordinary errors in the async multipart copy (a failed part or completion), and deliberately left cancellation out of scope.

Reproduction

Offline, on PR #1036 (0e75523), with the multipart requests mocked:

import asyncio, time
from types import SimpleNamespace
from unittest import mock
from pyathena.filesystem.s3 import S3FileSystem
from pyathena.filesystem.s3_async import AioS3FileSystem

async def main():
    fs = AioS3FileSystem(connection=mock.MagicMock(), max_workers=2, skip_instance_cache=True)
    sync_fs = fs._sync_fs
    sync_fs._call = mock.MagicMock(return_value={})
    sync_fs._create_multipart_upload = mock.MagicMock(return_value=SimpleNamespace(upload_id="u"))
    done = []
    def part(**kw):
        time.sleep(0.3); done.append(kw["part_number"])
        return SimpleNamespace(etag='"e"', part_number=kw["part_number"])
    sync_fs._upload_part_copy = mock.MagicMock(side_effect=part)
    sync_fs._complete_multipart_upload = mock.MagicMock()
    sync_fs._abort_multipart_upload = mock.MagicMock()
    task = asyncio.ensure_future(fs._copy_object_with_multipart_upload(
        bucket1="bucket", key1="src", size1=4 * S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE,
        bucket2="bucket", key2="dst", block_size=S3FileSystem.MULTIPART_UPLOAD_MIN_PART_SIZE,
        MetadataDirective="REPLACE", TaggingDirective="REPLACE", AnnotationDirective="EXCLUDE"))
    await asyncio.sleep(0.1)
    task.cancel()
    try:
        await task
    except asyncio.CancelledError:
        print("cancelled")
    print("abort called:", sync_fs._abort_multipart_upload.called, "parts finished:", sorted(done))
    await asyncio.sleep(0.5)
    print("parts finished after 0.5s:", sorted(done))

asyncio.run(main())

Output:

cancelled
abort called: False parts finished: []
parts finished after 0.5s: [1, 2]

The two running part copies complete after the cancellation has returned, and no AbortMultipartUpload is sent.

Expected

A cancelled async multipart copy waits for the part copies that are already running and then aborts the upload, as the sync path does for interrupts, and then re-raises the cancellation.

Environment

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions