| 497 | } |
| 498 | |
| 499 | void WriteBufferFromAzureBlobStorage::writePart(WriteBufferFromAzureBlobStorage::PartData && part_data) |
| 500 | { |
| 501 | const std::string & block_id = block_ids.emplace_back(getRandomASCIIString(64)); |
| 502 | auto worker_data = std::make_shared<std::tuple<std::string, WriteBufferFromAzureBlobStorage::PartData>>(block_id, std::move(part_data)); |
| 503 | |
| 504 | auto upload_worker = [this, worker_data] () |
| 505 | { |
| 506 | auto & data_size = std::get<1>(*worker_data).data_size; |
| 507 | auto & data_block_id = std::get<0>(*worker_data); |
| 508 | auto block_blob_client = blob_container_client->GetBlockBlobClient(blob_path); |
| 509 | |
| 510 | ProfileEvents::increment(ProfileEvents::AzureStageBlock); |
| 511 | if (blob_container_client->IsClientForDisk()) |
| 512 | ProfileEvents::increment(ProfileEvents::DiskAzureStageBlock); |
| 513 | |
| 514 | Azure::Core::IO::MemoryBodyStream memory_stream(reinterpret_cast<const uint8_t *>(std::get<1>(*worker_data).memory.data()), data_size); |
| 515 | |
| 516 | Stopwatch watch; |
| 517 | Int32 error_code = 0; |
| 518 | String error_message; |
| 519 | try |
| 520 | { |
| 521 | execWithRetry( |
| 522 | [&](size_t retry_attempt) |
| 523 | { |
| 524 | block_blob_client.StageBlock( |
| 525 | data_block_id, |
| 526 | memory_stream, |
| 527 | Azure::Storage::Blobs::StageBlockOptions{}, |
| 528 | azure_context.WithValue(PocoAzureHTTPClient::getSDKContextKeyForBufferRetry(), retry_attempt)); |
| 529 | }, |
| 530 | max_unexpected_write_error_retries, |
| 531 | data_size); |
| 532 | } |
| 533 | catch (const Azure::Core::RequestFailedException & e) |
| 534 | { |
| 535 | error_code = static_cast<Int32>(e.StatusCode); |
| 536 | error_message = e.Message; |
| 537 | if (blob_log) |
| 538 | blob_log->addEvent( |
| 539 | BlobStorageLogElement::EventType::MultiPartUploadWrite, |
| 540 | /* bucket */ container_for_logging, |
| 541 | /* remote_path */ blob_path, |
| 542 | /* local_path */ {}, |
| 543 | /* data_size */ data_size, |
| 544 | watch.elapsedMicroseconds(), |
| 545 | error_code, |
| 546 | error_message); |
| 547 | throw; |
| 548 | } |
| 549 | auto elapsed = watch.elapsedMicroseconds(); |
| 550 | |
| 551 | if (blob_log) |
| 552 | blob_log->addEvent( |
| 553 | BlobStorageLogElement::EventType::MultiPartUploadWrite, |
| 554 | /* bucket */ container_for_logging, |
| 555 | /* remote_path */ blob_path, |
| 556 | /* local_path */ {}, |
nothing calls this directly
no test coverage detected