diff --git a/pom.xml b/pom.xml index e1e32a9859..49af2e499c 100644 --- a/pom.xml +++ b/pom.xml @@ -114,7 +114,7 @@ 1.5.4 3.3.6 - 4.1.94.Final + 4.1.108.Final 1.2.1 2.1.6 @@ -956,6 +956,17 @@ pom import + + com.azure + azure-storage-blob + 12.29.0 + + + com.azure + azure-storage-common + 12.29.0 + + com.oceanbase apsara-audit-spring-boot-starter diff --git a/server/odc-core/src/main/java/com/oceanbase/odc/core/shared/constant/ConnectType.java b/server/odc-core/src/main/java/com/oceanbase/odc/core/shared/constant/ConnectType.java index 809a18dbb3..15064dfe8d 100644 --- a/server/odc-core/src/main/java/com/oceanbase/odc/core/shared/constant/ConnectType.java +++ b/server/odc-core/src/main/java/com/oceanbase/odc/core/shared/constant/ConnectType.java @@ -39,6 +39,8 @@ public enum ConnectType { OBS(DialectType.FILE_SYSTEM), COS(DialectType.FILE_SYSTEM), S3A(DialectType.FILE_SYSTEM), + // microsoft azure blob service + BLOB(DialectType.FILE_SYSTEM), UNKNOWN(DialectType.UNKNOWN), ; diff --git a/server/odc-service/pom.xml b/server/odc-service/pom.xml index 5d28d6077c..5f237653d7 100644 --- a/server/odc-service/pom.xml +++ b/server/odc-service/pom.xml @@ -67,6 +67,14 @@ io.micrometer micrometer-core + + com.azure + azure-storage-blob + + + com.azure + azure-storage-common + io.micrometer micrometer-registry-prometheus diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/connection/FileSystemConnectionTester.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/connection/FileSystemConnectionTester.java index 43d5c37be3..7fca6bf889 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/connection/FileSystemConnectionTester.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/connection/FileSystemConnectionTester.java @@ -32,6 +32,7 @@ import com.oceanbase.odc.service.cloud.model.CloudProvider; import com.oceanbase.odc.service.connection.model.ConnectionConfig; import com.oceanbase.odc.service.connection.model.ConnectionTestResult; +import com.oceanbase.odc.service.loaddata.model.ObjectStorageConfig; import com.oceanbase.odc.service.objectstorage.cloud.CloudResourceConfigurations; import com.oceanbase.odc.service.objectstorage.cloud.client.CloudClient; import com.oceanbase.odc.service.objectstorage.cloud.client.CloudException; @@ -70,7 +71,7 @@ public ConnectionTestResult test(@NonNull ConnectionConfig config) { storageConfig.setBucketName(uri.getAuthority()); storageConfig.setRegion(config.getRegion()); storageConfig.setCloudProvider(getCloudProvider(config.getType())); - storageConfig.setPublicEndpoint(getEndPointByRegion(config.getType(), config.getRegion())); + storageConfig.setPublicEndpoint(getEndPoint(config.getType(), config.getRegion(), config.getUsername())); try { CloudClient cloudClient = new CloudResourceConfigurations.CloudClientBuilder().generateCloudClient(storageConfig); @@ -123,7 +124,7 @@ private CloudProvider getCloudProvider(ConnectType type) { } } - private static String getEndPointByRegion(ConnectType type, String region) { + private static String getEndPoint(ConnectType type, String region, String userName) { switch (type) { case COS: return MessageFormat.format(COS_ENDPOINT_PATTERN, region); @@ -137,6 +138,8 @@ private static String getEndPointByRegion(ConnectType type, String region) { return MessageFormat.format(S3_ENDPOINT_CN_PATTERN, region); } return MessageFormat.format(S3_ENDPOINT_GLOBAL_PATTERN, region); + case BLOB: + return ObjectStorageConfig.concatEndpoint(CloudProvider.AZURE, userName); default: throw new IllegalArgumentException("regionToEndpoint is not applicable for storageType " + type); } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/model/ObjectStorageConfig.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/model/ObjectStorageConfig.java index 8d8fb6d111..e07fc37e79 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/model/ObjectStorageConfig.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/model/ObjectStorageConfig.java @@ -189,6 +189,7 @@ public String getEndpoint() { // while azure use the account name as one component in endpoint. // GCS need neither region nor access key to concat the endpoint. if (cloudProvider == CloudProvider.AZURE) { + // for azure, accessKey is account name, secretKey is the secret key this.endpoint = concatEndpoint(cloudProvider, accessKey); } else if (cloudProvider == CloudProvider.GOOGLE_CLOUD) { this.endpoint = GCS_ENDPOINT; @@ -217,7 +218,7 @@ public ObjectStorageConfiguration toObjectStorageConfiguration() { return config; } - private static String concatEndpoint(CloudProvider cloudProvider, String component) { + public static String concatEndpoint(CloudProvider cloudProvider, String component) { switch (cloudProvider) { case TENCENT_CLOUD: return MessageFormat.format(COS_ENDPOINT_PATTERN, component); diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/util/CloudProviderUtil.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/util/CloudProviderUtil.java index 5a97ad3fc9..0b54aacbbe 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/util/CloudProviderUtil.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/loaddata/util/CloudProviderUtil.java @@ -46,24 +46,4 @@ public static CloudProvider fromScheme(String scheme) { throw new IllegalArgumentException("Unsupported scheme: " + scheme); } } - - public static String getObjectUri(CloudProvider provider, String bucket, String objectName) { - switch (provider) { - case ALIBABA_CLOUD: - return "oss://" + bucket + "/" + objectName; - case TENCENT_CLOUD: - return "cos://" + bucket + "/" + objectName; - case HUAWEI_CLOUD: - return "obs://" + bucket + "/" + objectName; - case AWS: - return "s3://" + bucket + "/" + objectName; - case AZURE: - return "azblob://" + bucket + "/" + objectName; - case GOOGLE_CLOUD: - return "gcs://" + bucket + "/" + objectName; - default: - throw new IllegalArgumentException("Unsupported cloud provider: " + provider); - } - } - } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/client/CloudObjectStorageClient.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/client/CloudObjectStorageClient.java index eb3ca19244..a2333e69f0 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/client/CloudObjectStorageClient.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/client/CloudObjectStorageClient.java @@ -290,7 +290,7 @@ private CompleteMultipartUploadResult multiPartUpload(@NotBlank String objectNam } } CompleteMultipartUploadRequest completeMultipartUploadRequest = - new CompleteMultipartUploadRequest(bucketName, objectName, uploadId, partTags); + new CompleteMultipartUploadRequest(bucketName, objectName, uploadId, partTags, metadata); CompleteMultipartUploadResult completeMultipartUploadResult = internalEndpointCloudObjectStorage.completeMultipartUpload(completeMultipartUploadRequest); log.info("Complete multipart upload, result={}", completeMultipartUploadResult); diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/CloudResourceConfigurations.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/CloudResourceConfigurations.java index c03a4fd259..eed4952a63 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/CloudResourceConfigurations.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/CloudResourceConfigurations.java @@ -39,11 +39,15 @@ import com.amazonaws.services.s3.AmazonS3ClientBuilder; import com.amazonaws.services.securitytoken.AWSSecurityTokenService; import com.amazonaws.services.securitytoken.AWSSecurityTokenServiceClientBuilder; +import com.azure.storage.blob.BlobServiceClient; +import com.azure.storage.blob.BlobServiceClientBuilder; +import com.azure.storage.common.StorageSharedKeyCredential; import com.oceanbase.odc.common.util.StringUtils; import com.oceanbase.odc.core.shared.PreConditions; import com.oceanbase.odc.service.cloud.model.CloudProvider; import com.oceanbase.odc.service.objectstorage.cloud.client.AlibabaCloudClient; import com.oceanbase.odc.service.objectstorage.cloud.client.AmazonCloudClient; +import com.oceanbase.odc.service.objectstorage.cloud.client.AzureCloudClient; import com.oceanbase.odc.service.objectstorage.cloud.client.CloudClient; import com.oceanbase.odc.service.objectstorage.cloud.client.GoogleCloudClient; import com.oceanbase.odc.service.objectstorage.cloud.client.NullCloudClient; @@ -110,6 +114,8 @@ public CloudClient generateCloudClient(ObjectStorageConfiguration configuration, return createAmazonCloudClient(configuration); case GOOGLE_CLOUD: return createGoogleCloudClient(configuration); + case AZURE: + return createAzureCloudClient(configuration); default: return new NullCloudClient(); } @@ -198,6 +204,18 @@ static GoogleCloudClient createGoogleCloudClient(ObjectStorageConfiguration conf return new GoogleCloudClient(s3, sts, roleSessionName, roleArn); } + static AzureCloudClient createAzureCloudClient(ObjectStorageConfiguration configuration) { + // for azure, accessKey is account name, secretKey is the secret key + String accountName = configuration.getAccessKeyId(); + String accountKey = configuration.getAccessKeySecret(); + StorageSharedKeyCredential credential = new StorageSharedKeyCredential(accountName, accountKey); + BlobServiceClient serviceClient = new BlobServiceClientBuilder() + .endpoint(configuration.getPublicEndpoint()) + .credential(credential) + .buildClient(); + return new AzureCloudClient(serviceClient, configuration.getRegion()); + } + private static OSS getInternalOss(ObjectStorageConfiguration objectStorageConfiguration) { return createOssClient( objectStorageConfiguration.getAccessKeyId(), diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/client/AzureCloudClient.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/client/AzureCloudClient.java new file mode 100644 index 0000000000..39242c904a --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/client/AzureCloudClient.java @@ -0,0 +1,506 @@ +/* + * Copyright (c) 2023 OceanBase. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.oceanbase.odc.service.objectstorage.cloud.client; + +import java.io.ByteArrayInputStream; +import java.io.File; +import java.io.IOException; +import java.io.InputStream; +import java.net.URL; +import java.nio.charset.StandardCharsets; +import java.time.Instant; +import java.time.OffsetDateTime; +import java.time.ZoneId; +import java.time.temporal.ChronoUnit; +import java.util.ArrayList; +import java.util.Base64; +import java.util.Collection; +import java.util.Date; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.function.Consumer; +import java.util.stream.Collectors; + +import org.apache.commons.lang3.StringUtils; + +import com.azure.core.http.rest.PagedIterable; +import com.azure.core.http.rest.Response; +import com.azure.core.util.ProgressListener; +import com.azure.storage.blob.BlobClient; +import com.azure.storage.blob.BlobContainerClient; +import com.azure.storage.blob.BlobServiceClient; +import com.azure.storage.blob.models.BlobHttpHeaders; +import com.azure.storage.blob.models.BlobItem; +import com.azure.storage.blob.models.BlockBlobItem; +import com.azure.storage.blob.models.ListBlobsOptions; +import com.azure.storage.blob.models.ParallelTransferOptions; +import com.azure.storage.blob.options.BlobParallelUploadOptions; +import com.azure.storage.blob.options.BlobUploadFromFileOptions; +import com.azure.storage.blob.sas.BlobSasPermission; +import com.azure.storage.blob.sas.BlobServiceSasSignatureValues; +import com.azure.storage.blob.specialized.BlobInputStream; +import com.azure.storage.blob.specialized.BlockBlobClient; +import com.oceanbase.odc.service.objectstorage.cloud.model.CloudObjectStorageConstants; +import com.oceanbase.odc.service.objectstorage.cloud.model.CompleteMultipartUploadRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.CompleteMultipartUploadResult; +import com.oceanbase.odc.service.objectstorage.cloud.model.CopyObjectResult; +import com.oceanbase.odc.service.objectstorage.cloud.model.DeleteObjectRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.DeleteObjectsRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.DeleteObjectsResult; +import com.oceanbase.odc.service.objectstorage.cloud.model.GetObjectRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.InitiateMultipartUploadRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.InitiateMultipartUploadResult; +import com.oceanbase.odc.service.objectstorage.cloud.model.ObjectMetadata; +import com.oceanbase.odc.service.objectstorage.cloud.model.ObjectSummary; +import com.oceanbase.odc.service.objectstorage.cloud.model.ObjectTagging; +import com.oceanbase.odc.service.objectstorage.cloud.model.PartETag; +import com.oceanbase.odc.service.objectstorage.cloud.model.PutObjectResult; +import com.oceanbase.odc.service.objectstorage.cloud.model.StorageObject; +import com.oceanbase.odc.service.objectstorage.cloud.model.UploadObjectTemporaryCredential; +import com.oceanbase.odc.service.objectstorage.cloud.model.UploadPartRequest; +import com.oceanbase.odc.service.objectstorage.cloud.model.UploadPartResult; + +import cn.hutool.core.util.HexUtil; +import lombok.extern.slf4j.Slf4j; + +/** + * @author longpeng.zlp + * @date 2025/4/24 15:05 + */ +@Slf4j +public class AzureCloudClient implements CloudClient { + private final BlobServiceClient blobServiceClient; + private final String region; + + public AzureCloudClient(BlobServiceClient blobServiceClient, String region) { + this.blobServiceClient = blobServiceClient; + this.region = region; + } + + @Override + public boolean supported() { + return true; + } + + // bucket location is a attribute for storage account with is parent concept of blob service + // azure provide big area concept as east asia, japan asia, pacific australia eg. + // so we return null as we don't care about the location + @Override + public String getBucketLocation(String bucketName) throws CloudException { + return region; + } + + @Override + public boolean doesBucketExist(String bucketName) throws CloudException { + return callAzureMethod("does bucket exist " + bucketName, () -> { + BlobContainerClient blobContainerClient = blobServiceClient.getBlobContainerClient(bucketName); + return blobContainerClient.exists(); + }); + } + + @Override + public InitiateMultipartUploadResult initiateMultipartUpload(InitiateMultipartUploadRequest request) + throws CloudException { + InitiateMultipartUploadResult ret = new InitiateMultipartUploadResult(); + ret.setBucketName(request.getBucketName()); + ret.setKey(request.getKey()); + ret.setUploadId("AzurePartUpload"); + return ret; + } + + @Override + public UploadPartResult uploadPart(UploadPartRequest request) throws CloudException { + return callAzureMethod("upload part object " + request.getBucketName() + "." + request.getKey(), () -> { + BlobClient blobClient = getBlobClient(request.getBucketName(), request.getKey()); + BlockBlobClient client = blobClient.getBlockBlobClient(); + String blockID = + Base64.getEncoder().encodeToString(UUID.randomUUID().toString().getBytes(StandardCharsets.UTF_8)); + InputStream bufferedInputStream = + castToRemarkableStream(request.getInputStream(), (int) request.getPartSize()); + client.stageBlock(blockID, bufferedInputStream, request.getPartSize()); + UploadPartResult result = new UploadPartResult(); + result.setPartNumber(result.getPartNumber()); + result.setPartSize(request.getPartSize()); + PartETag partETag = new PartETag(); + partETag.setETag(blockID); + result.setPartETag(partETag); + return result; + }); + } + + private InputStream castToRemarkableStream(InputStream inputStream, int size) throws IOException { + byte[] buffer = new byte[size]; + int remain = size; + int offset = 0; + int readByte = 0; + while (0 != remain && (readByte = inputStream.read(buffer, offset, remain)) != -1) { + remain -= readByte; + offset += readByte; + } + return new ByteArrayInputStream(buffer); + } + + @Override + public CompleteMultipartUploadResult completeMultipartUpload(CompleteMultipartUploadRequest request) + throws CloudException { + return callAzureMethod("complete upload part object " + request.getBucketName() + "." + request.getKey(), + () -> { + BlobClient blobClient = getBlobClient(request.getBucketName(), request.getKey()); + BlockBlobClient client = blobClient.getBlockBlobClient(); + List blockIDs = + request.getPartETags().stream().map(PartETag::getETag).collect(Collectors.toList()); + client.commitBlockList(blockIDs, true); + fillOptions(blobClient::setMetadata, blobClient::setTags, blobClient::setHttpHeaders, + request.getObjectMetadata()); + CompleteMultipartUploadResult ret = new CompleteMultipartUploadResult(); + ret.setBucketName(request.getBucketName()); + ret.setKey(request.getKey()); + return ret; + }); + } + + // if file has exists, we use cover semantic + @Override + public PutObjectResult putObject(String bucketName, String key, File file, ObjectMetadata metadata) + throws CloudException { + return callAzureMethod("put object " + bucketName + "." + file.getName(), () -> { + BlobClient client = getBlobClient(bucketName, key); + if (client.exists()) { + log.info("{}.{} has exists, cover it", bucketName, key); + } + BlobUploadFromFileOptions options = buildUploadFromFileOptions(file, metadata); + Response response = client.uploadFromFileWithResponse(options, null, null); + return fromUploadResponse(response); + }); + } + + @Override + public PutObjectResult putObject(String bucketName, String key, InputStream in, ObjectMetadata metadata) + throws CloudException { + return callAzureMethod("put object " + bucketName + "." + key, () -> { + BlobClient client = getBlobClient(bucketName, key); + if (client.exists()) { + log.info("{}.{} has exists, cover it", bucketName, key); + } + BlobParallelUploadOptions options = buildUploadInputStreamOptions(in, metadata); + Response response = client.uploadWithResponse(options, null, null); + return fromUploadResponse(response); + }); + } + + protected PutObjectResult fromUploadResponse(Response response) { + BlockBlobItem blockBlobItem = response.getValue(); + PutObjectResult ret = new PutObjectResult(); + ret.setETag(blockBlobItem.getETag()); + ret.setVersionId(blockBlobItem.getVersionId()); + ret.setRequestId(response.getRequest().getUrl().toString()); + return ret; + } + + protected BlobUploadFromFileOptions buildUploadFromFileOptions(File file, ObjectMetadata metadata) { + BlobUploadFromFileOptions blobUploadFromFileOptions = new BlobUploadFromFileOptions(file.getAbsolutePath()); + blobUploadFromFileOptions.setParallelTransferOptions(buildParallelTransferOptions()); + fillOptions(blobUploadFromFileOptions::setMetadata, blobUploadFromFileOptions::setTags, + blobUploadFromFileOptions::setHeaders, metadata); + return blobUploadFromFileOptions; + } + + protected BlobParallelUploadOptions buildUploadInputStreamOptions(InputStream stream, ObjectMetadata metadata) { + BlobParallelUploadOptions blobParallelUploadOptions = new BlobParallelUploadOptions(stream); + blobParallelUploadOptions.setParallelTransferOptions(buildParallelTransferOptions()); + fillOptions(blobParallelUploadOptions::setMetadata, blobParallelUploadOptions::setTags, + blobParallelUploadOptions::setHeaders, metadata); + return blobParallelUploadOptions; + } + + protected void fillOptions(Consumer> metadataConsumer, + Consumer> tagConsumer, + Consumer headersConsumer, + ObjectMetadata metadata) { + if (null == metadata) { + return; + } + if (null != metadata.getUserMetadata()) { + metadataConsumer.accept(metadata.getUserMetadata()); + } + if (metadata.hasTag()) { + tagConsumer.accept(castTag(metadata)); + } + BlobHttpHeaders httpHeaders = new BlobHttpHeaders(); + boolean headerValid = false; + if (null != metadata.getContentType()) { + httpHeaders.setContentType(metadata.getContentType()); + headerValid = true; + } + if (null != metadata.getContentMD5()) { + byte[] bytes = HexUtil.decodeHex(metadata.getContentMD5()); + httpHeaders.setContentMd5(bytes); + headerValid = true; + } + if (headerValid) { + headersConsumer.accept(httpHeaders); + } + } + + @Override + public CopyObjectResult copyObject(String bucketName, String from, String to) throws CloudException { + return callAzureMethod("copy object " + bucketName + "." + from + " to " + to, () -> { + BlobClient destClient = getBlobClient(bucketName, to); + // destClient.copyFromUrl(srcClient.getBlobUrl()); + String originURL = generatePresignedUrl(bucketName, from, + Date.from(Instant.ofEpochMilli(System.currentTimeMillis() + 3600 * 1000))).toString(); + // tag will not be copied + destClient.copyFromUrl(originURL); + CopyObjectResult result = new CopyObjectResult(); + destClient = getBlobClient(bucketName, to); + result.setETag(destClient.getProperties().getETag()); + result.setVersionId(destClient.getVersionId()); + result.setLastModifyTime(new Date(destClient.getProperties().getLastModified().toEpochSecond())); + return result; + }); + } + + @Override + public DeleteObjectsResult deleteObjects(DeleteObjectsRequest request) throws CloudException { + return callAzureMethod( + "delete objects " + request.getBucketName() + "." + StringUtils.join(request.getKeys(), ","), () -> { + List ret = new ArrayList<>(); + BlobContainerClient blobContainerClient = + blobServiceClient.getBlobContainerClient(request.getBucketName()); + for (String key : request.getKeys()) { + BlobClient blobClient = blobContainerClient.getBlobClient(key); + if (blobClient.deleteIfExists()) { + ret.add(key); + } + } + DeleteObjectsResult result = new DeleteObjectsResult(); + result.setDeletedObjects(ret); + return result; + }); + } + + @Override + public String deleteObject(DeleteObjectRequest request) throws CloudException { + return callAzureMethod("delete object " + request.getBucketName() + "." + request.getKey(), () -> { + BlobClient blobClient = getBlobClient(request.getBucketName(), request.getKey()); + blobClient.deleteIfExists(); + return request.getKey(); + }); + } + + @Override + public boolean doesObjectExist(String bucketName, String key) throws CloudException { + return callAzureMethod("does object exist " + bucketName + "." + key, () -> { + BlobClient blobClient = getBlobClient(bucketName, key); + return blobClient.exists(); + }); + } + + @Override + public StorageObject getObject(String bucketName, String key) throws CloudException { + return callAzureMethod("get object " + bucketName + "." + key, () -> { + BlobClient blobClient = getBlobClient(bucketName, key); + BlobInputStream inputStream = blobClient.openInputStream(); + StorageObject ret = new StorageObject() { + @Override + public InputStream getObjectContent() { + return inputStream; + } + + @Override + public InputStream getAbortableContent() { + return inputStream; + } + }; + ret.setBucketName(bucketName); + ret.setKey(key); + ret.setMetadata(getObjectMetadata(blobClient)); + return ret; + }); + } + + @Override + public ObjectMetadata getObject(GetObjectRequest request, File file) throws CloudException { + return callAzureMethod("get object " + request.getKey() + " to " + file.getName(), () -> { + BlobClient blobClient = getBlobClient(request.getBucketName(), request.getKey()); + blobClient.downloadToFile(file.getAbsolutePath()); + return getObjectMetadata(blobClient); + }); + } + + @Override + public ObjectMetadata getObjectMetadata(String bucketName, String key) throws CloudException { + return callAzureMethod("get object metadata " + bucketName + "." + key, () -> { + BlobClient blobClient = getBlobClient(bucketName, key); + if (!blobClient.exists()) { + return null; + } + return getObjectMetadata(blobClient); + }); + } + + protected ObjectMetadata getObjectMetadata(BlobClient client) { + ObjectMetadata ret = new ObjectMetadata(); + ret.setContentType(client.getProperties().getContentType()); + ret.setETag(client.getProperties().getETag()); + ret.setUserMetadata(client.getProperties().getMetadata()); + ret.setContentLength(client.getProperties().getBlobSize()); + if (null != client.getProperties().getExpiresOn()) { + ret.setExpirationTime(new Date(client.getProperties().getExpiresOn().toEpochSecond())); + } + if (null != client.getTags()) { + ObjectTagging objectTagging = new ObjectTagging(); + client.getTags().forEach((k, v) -> { + objectTagging.withTag(k, v); + }); + ret.setTagging(objectTagging); + } + return ret; + } + + @Override + public URL generatePresignedUrl(String bucketName, String key, Date expiration) throws CloudException { + return generatePresignedUrlWithCustomFileName(bucketName, key, expiration, null); + } + + @Override + public URL generatePresignedUrlWithCustomFileName(String bucketName, String key, Date expiration, + String customFileName) throws CloudException { + return callAzureMethod("generate pregisgned url for " + bucketName + "." + key, () -> { + BlobClient blobClient = getBlobClient(bucketName, key); + BlobSasPermission permission = new BlobSasPermission(); + permission.setReadPermission(true); + permission.setTagsPermission(true); + BlobServiceSasSignatureValues sasValues = new BlobServiceSasSignatureValues() + .setPermissions(permission); + if (null != expiration) { + sasValues.setExpiryTime(OffsetDateTime.ofInstant(expiration.toInstant(), ZoneId.systemDefault())); + } + if (!StringUtils.isEmpty(customFileName)) { + sasValues.setContentDisposition(String.format("attachment;filename=%s", customFileName)); + } + String sasToken = blobClient.generateSas(sasValues); + String url = blobClient.getBlobUrl() + "?" + sasToken; + return new URL(url); + }); + } + + @Override + public URL generatePresignedPutUrl(String bucketName, String key, Date expiration) throws CloudException { + return callAzureMethod("generate pregisgned url for " + bucketName + "." + key, () -> { + BlobClient blobClient = getBlobClient(bucketName, key); + BlobSasPermission permission = new BlobSasPermission(); + permission.setWritePermission(true); + permission.setCreatePermission(true); + permission.setTagsPermission(true); + BlobServiceSasSignatureValues sasValues = new BlobServiceSasSignatureValues() + .setPermissions(permission); + if (null != expiration) { + sasValues.setExpiryTime(OffsetDateTime.ofInstant(expiration.toInstant(), ZoneId.systemDefault())); + } + String sasToken = blobClient.generateSas(sasValues); + String url = blobClient.getBlobUrl() + "?" + sasToken; + return new URL(url); + }); + } + + @Override + public List list(String bucketName, String prefix) throws CloudException { + return callAzureMethod("list " + bucketName + "." + prefix, () -> { + BlobContainerClient blobContainerClient = blobServiceClient.getBlobContainerClient(bucketName); + ListBlobsOptions listBlobsOptions = new ListBlobsOptions(); + listBlobsOptions.setPrefix(prefix); + PagedIterable blobItems = blobContainerClient.listBlobs(listBlobsOptions, null); + List ret = new ArrayList<>(); + for (BlobItem item : blobItems) { + ObjectSummary tmp = new ObjectSummary(); + tmp.setBucketName(bucketName); + tmp.setKey(item.getName()); + tmp.setETag(item.getProperties().getETag()); + tmp.setSize(item.getProperties().getContentLength()); + tmp.setType(item.getProperties().getContentType()); + tmp.setLastModified(new Date(item.getProperties().getLastModified().toEpochSecond())); + tmp.setStorageClass(item.getProperties().getAccessTier().getValue()); + ret.add(tmp); + } + return ret; + }); + } + + @Override + public UploadObjectTemporaryCredential generateTempCredential(String bucketName, String objectName, + Long durationSeconds) throws CloudException { + return callAzureMethod("generate temp credential " + bucketName + "." + objectName, () -> { + BlobClient blobClient = getBlobClient(bucketName, objectName); + BlobSasPermission permission = BlobSasPermission.parse("acdelimrtwxy"); + BlobServiceSasSignatureValues sasValues = new BlobServiceSasSignatureValues() + .setPermissions(permission); + if (null != durationSeconds) { + sasValues.setExpiryTime(OffsetDateTime + .ofInstant(Instant.now().plus(durationSeconds, ChronoUnit.SECONDS), ZoneId.systemDefault())); + } + String sasToken = blobClient.generateSas(sasValues); + UploadObjectTemporaryCredential ret = new UploadObjectTemporaryCredential(); + ret.setSecurityToken(sasToken); + return ret; + }); + } + + protected BlobClient getBlobClient(String bucketName, String key) { + BlobContainerClient blobContainerClient = blobServiceClient.getBlobContainerClient(bucketName); + return blobContainerClient.getBlobClient(key); + } + + protected ParallelTransferOptions buildParallelTransferOptions() { + ParallelTransferOptions ret = new ParallelTransferOptions(); + ret.setBlockSizeLong(CloudObjectStorageConstants.MIN_PART_SIZE); + ret.setMaxSingleUploadSizeLong(CloudObjectStorageConstants.CRITICAL_FILE_SIZE_IN_MB); + ret.setMaxConcurrency(2); + ret.setProgressListener(new ProgressListener() { + @Override + public void handleProgress(long l) { + // print progress + } + }); + return ret; + } + + protected Map castTag(ObjectMetadata objectMetadata) { + if (null == objectMetadata || !objectMetadata.hasTag()) { + return null; + } + Collection tagList = objectMetadata.getTagging().getTagSet(); + Map tapMap = new HashMap<>(); + for (ObjectTagging.Tag tag : tagList) { + tapMap.put(tag.getKey(), tag.getValue()); + } + return tapMap; + } + + protected T callAzureMethod(String operation, AzureOperator supplier) { + try { + return supplier.getResult(); + } catch (Exception ex) { + throw new CloudException(operation + " failed", ex); + } + } + + private interface AzureOperator { + T getResult() throws Exception; + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/CompleteMultipartUploadRequest.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/CompleteMultipartUploadRequest.java index d5f99dfa05..36eb05dbd7 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/CompleteMultipartUploadRequest.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/CompleteMultipartUploadRequest.java @@ -27,10 +27,13 @@ public class CompleteMultipartUploadRequest extends GenericRequest { private String uploadId; private List partETags; + private ObjectMetadata objectMetadata; - public CompleteMultipartUploadRequest(String bucketName, String key, String uploadId, List partETags) { + public CompleteMultipartUploadRequest(String bucketName, String key, String uploadId, List partETags, + ObjectMetadata objectMetadata) { super(bucketName, key); this.uploadId = uploadId; this.partETags = partETags; + this.objectMetadata = objectMetadata; } } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/ObjectStorageConfiguration.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/ObjectStorageConfiguration.java index 3fe5ce3e3d..b4df570fc4 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/ObjectStorageConfiguration.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/cloud/model/ObjectStorageConfiguration.java @@ -23,6 +23,7 @@ import com.oceanbase.odc.common.util.StringUtils; import com.oceanbase.odc.core.shared.exception.UnexpectedException; import com.oceanbase.odc.service.cloud.model.CloudProvider; +import com.oceanbase.odc.service.loaddata.model.ObjectStorageConfig; import lombok.Data; @@ -50,20 +51,25 @@ public class ObjectStorageConfiguration { * for aws s3, if endpoint not set, get by region */ public String getPublicEndpoint() { - if (StringUtils.isBlank(publicEndpoint) && StringUtils.isNotBlank(region)) { - if (CloudProvider.AWS == cloudProvider) { - return RegionUtils.getRegion(region).getServiceEndpoint("s3"); - } else if (CloudProvider.ALIBABA_CLOUD == cloudProvider) { - try { - ResolveEndpointRequest request = new ResolveEndpointRequest(region, "oss", null, null); - EndpointResolver endpointResolver = new LocalConfigRegionalEndpointResolver(); - return endpointResolver.resolve(request); - } catch (ClientException e) { - throw new UnexpectedException("getProfile failed with region=" + region, e); - } + // return if set + if (!StringUtils.isBlank(publicEndpoint)) { + return publicEndpoint; + } + if (CloudProvider.AWS == cloudProvider && StringUtils.isNotBlank(region)) { + return RegionUtils.getRegion(region).getServiceEndpoint("s3"); + } else if (CloudProvider.ALIBABA_CLOUD == cloudProvider && StringUtils.isNotBlank(region)) { + try { + ResolveEndpointRequest request = new ResolveEndpointRequest(region, "oss", null, null); + EndpointResolver endpointResolver = new LocalConfigRegionalEndpointResolver(); + return endpointResolver.resolve(request); + } catch (ClientException e) { + throw new UnexpectedException("getProfile failed with region=" + region, e); } + } else if (CloudProvider.AZURE == cloudProvider) { + // for azure compute as it's endpoint + return ObjectStorageConfig.concatEndpoint(cloudProvider, accessKeySecret); } - return this.publicEndpoint; + return publicEndpoint; } public String getInternalEndpoint() { diff --git a/server/odc-service/src/test/java/com/oceanbase/odc/service/objectstorage/AzureCloudClientTest.java b/server/odc-service/src/test/java/com/oceanbase/odc/service/objectstorage/AzureCloudClientTest.java new file mode 100644 index 0000000000..f4dd1e9736 --- /dev/null +++ b/server/odc-service/src/test/java/com/oceanbase/odc/service/objectstorage/AzureCloudClientTest.java @@ -0,0 +1,59 @@ +/* + * Copyright (c) 2023 OceanBase. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.oceanbase.odc.service.objectstorage; + +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mockito; + +import com.azure.storage.blob.BlobContainerClient; +import com.azure.storage.blob.BlobServiceClient; +import com.oceanbase.odc.service.objectstorage.cloud.client.AzureCloudClient; + +/** + * @author longpeng.zlp + * @date 2025/5/12 11:46 + */ +public class AzureCloudClientTest { + + private AzureCloudClient azureCloudClient; + private BlobServiceClient serviceClient; + + @Before + public void init() { + serviceClient = Mockito.mock(BlobServiceClient.class); + azureCloudClient = new AzureCloudClient(serviceClient, "ss"); + } + + @Test + public void testBucketExits() { + BlobContainerClient existClient = Mockito.mock(BlobContainerClient.class); + Mockito.when(existClient.exists()).thenReturn(true); + BlobContainerClient notExistClient = Mockito.mock(BlobContainerClient.class); + Mockito.when(notExistClient.exists()).thenReturn(false); + + Mockito.when(serviceClient.getBlobContainerClient("odccontainer")).thenReturn(existClient); + Mockito.when(serviceClient.getBlobContainerClient("ssssss")).thenReturn(notExistClient); + Assert.assertTrue(azureCloudClient.doesBucketExist("odccontainer")); + Assert.assertFalse(azureCloudClient.doesBucketExist("ssssss")); + } + + @Test + public void testGetLocation() { + Assert.assertEquals("ss", azureCloudClient.getBucketLocation("odccontainer")); + } +}