sijie closed pull request #1746: PIP-17: impl offload() for
S3ManagedLedgerOffloader
URL: https://github.com/apache/incubator-pulsar/pull/1746
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
diff --git a/conf/broker.conf b/conf/broker.conf
index e4318109cb..cf9faa8cb2 100644
--- a/conf/broker.conf
+++ b/conf/broker.conf
@@ -492,6 +492,9 @@ s3ManagedLedgerOffloadBucket=
# For Amazon S3 ledger offload, Alternative endpoint to connect to (useful for
testing)
s3ManagedLedgerOffloadServiceEndpoint=
+# For Amazon S3 ledger offload, Max block size in bytes.
+s3ManagedLedgerOffloadMaxBlockSizeInBytes=67108864
+
### --- Deprecated config variables --- ###
# Deprecated. Use configurationStoreServers
diff --git
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
index 43d929d58e..74363cc142 100644
---
a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
+++
b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java
@@ -484,6 +484,9 @@
// For Amazon S3 ledger offload, Alternative endpoint to connect to
(useful for testing)
private String s3ManagedLedgerOffloadServiceEndpoint = null;
+ // For Amazon S3 ledger offload, Max block size in bytes.
+ private int s3ManagedLedgerOffloadMaxBlockSizeInBytes = 64 * 1024 * 1024;
+
public String getZookeeperServers() {
return zookeeperServers;
}
@@ -1682,4 +1685,13 @@ public void
setS3ManagedLedgerOffloadServiceEndpoint(String endpoint) {
public String getS3ManagedLedgerOffloadServiceEndpoint() {
return this.s3ManagedLedgerOffloadServiceEndpoint;
}
+
+ public void setS3ManagedLedgerOffloadMaxBlockSizeInBytes(int
blockSizeInBytes) {
+ this.s3ManagedLedgerOffloadMaxBlockSizeInBytes = blockSizeInBytes;
+ }
+
+ public int getS3ManagedLedgerOffloadMaxBlockSizeInBytes() {
+ return this.s3ManagedLedgerOffloadMaxBlockSizeInBytes;
+ }
+
}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
index 163b79b23e..76920debe3 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
@@ -18,27 +18,45 @@
*/
package org.apache.pulsar.broker.s3offload;
+import com.amazonaws.AmazonServiceException;
+import com.amazonaws.SdkClientException;
import com.amazonaws.client.builder.AwsClientBuilder.EndpointConfiguration;
import com.amazonaws.services.s3.AmazonS3;
import com.amazonaws.services.s3.AmazonS3ClientBuilder;
-
+import com.amazonaws.services.s3.model.AbortMultipartUploadRequest;
+import com.amazonaws.services.s3.model.CompleteMultipartUploadRequest;
+import com.amazonaws.services.s3.model.InitiateMultipartUploadRequest;
+import com.amazonaws.services.s3.model.InitiateMultipartUploadResult;
+import com.amazonaws.services.s3.model.ObjectMetadata;
+import com.amazonaws.services.s3.model.PartETag;
+import com.amazonaws.services.s3.model.PutObjectRequest;
+import com.amazonaws.services.s3.model.UploadPartRequest;
+import com.amazonaws.services.s3.model.UploadPartResult;
import com.google.common.base.Strings;
-
+import java.io.InputStream;
+import java.util.LinkedList;
+import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ScheduledExecutorService;
-
import org.apache.bookkeeper.client.api.ReadHandle;
import org.apache.bookkeeper.mledger.LedgerOffloader;
import org.apache.pulsar.broker.PulsarServerException;
import org.apache.pulsar.broker.ServiceConfiguration;
+import
org.apache.pulsar.broker.s3offload.impl.BlockAwareSegmentInputStreamImpl;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
public class S3ManagedLedgerOffloader implements LedgerOffloader {
+ private static final Logger log =
LoggerFactory.getLogger(S3ManagedLedgerOffloader.class);
+
public static final String DRIVER_NAME = "S3";
private final ScheduledExecutorService scheduler;
private final AmazonS3 s3client;
private final String bucket;
+ // max block size for each data block.
+ private int maxBlockSize;
public static S3ManagedLedgerOffloader create(ServiceConfiguration conf,
ScheduledExecutorService
scheduler)
@@ -46,6 +64,8 @@ public static S3ManagedLedgerOffloader
create(ServiceConfiguration conf,
String region = conf.getS3ManagedLedgerOffloadRegion();
String bucket = conf.getS3ManagedLedgerOffloadBucket();
String endpoint = conf.getS3ManagedLedgerOffloadServiceEndpoint();
+ int maxBlockSize = conf.getS3ManagedLedgerOffloadMaxBlockSizeInBytes();
+
if (Strings.isNullOrEmpty(region)) {
throw new PulsarServerException("s3ManagedLedgerOffloadRegion
cannot be empty is s3 offload enabled");
}
@@ -60,28 +80,110 @@ public static S3ManagedLedgerOffloader
create(ServiceConfiguration conf,
} else {
builder.setRegion(region);
}
- return new S3ManagedLedgerOffloader(builder.build(), bucket,
scheduler);
+ return new S3ManagedLedgerOffloader(builder.build(), bucket,
scheduler, maxBlockSize);
}
- S3ManagedLedgerOffloader(AmazonS3 s3client, String bucket,
ScheduledExecutorService scheduler) {
+ S3ManagedLedgerOffloader(AmazonS3 s3client, String bucket,
ScheduledExecutorService scheduler, int maxBlockSize) {
this.s3client = s3client;
this.bucket = bucket;
this.scheduler = scheduler;
+ this.maxBlockSize = maxBlockSize;
+ }
+
+ static String dataBlockOffloadKey(ReadHandle readHandle, UUID uuid) {
+ return String.format("ledger-%d-%s", readHandle.getId(),
uuid.toString());
+ }
+
+ static String indexBlockOffloadKey(ReadHandle readHandle, UUID uuid) {
+ return String.format("ledger-%d-%s-index", readHandle.getId(),
uuid.toString());
}
+ // upload DataBlock to s3 using MultiPartUpload, and indexBlock in a new
Block,
@Override
- public CompletableFuture<Void> offload(ReadHandle ledger,
- UUID uid,
+ public CompletableFuture<Void> offload(ReadHandle readHandle,
+ UUID uuid,
Map<String, String> extraMetadata) {
CompletableFuture<Void> promise = new CompletableFuture<>();
scheduler.submit(() -> {
- try {
- s3client.putObject(bucket, uid.toString(), uid.toString());
- promise.complete(null);
- } catch (Throwable t) {
- promise.completeExceptionally(t);
+ OffloadIndexBlockBuilder indexBuilder =
OffloadIndexBlockBuilder.create()
+ .withMetadata(readHandle.getLedgerMetadata());
+ String dataBlockKey = dataBlockOffloadKey(readHandle, uuid);
+ String indexBlockKey = indexBlockOffloadKey(readHandle, uuid);
+ InitiateMultipartUploadRequest dataBlockReq = new
InitiateMultipartUploadRequest(bucket, dataBlockKey);
+ InitiateMultipartUploadResult dataBlockRes = null;
+
+ // init multi part upload for data block.
+ try {
+ dataBlockRes = s3client.initiateMultipartUpload(dataBlockReq);
+ } catch (Throwable t) {
+ promise.completeExceptionally(t);
+ return;
+ }
+
+ // start multi part upload for data block.
+ try {
+ long startEntry = 0;
+ int partId = 1;
+ long entryBytesWritten = 0;
+ List<PartETag> etags = new LinkedList<>();
+ while (startEntry <= readHandle.getLastAddConfirmed()) {
+ int blockSize = BlockAwareSegmentInputStreamImpl
+ .calculateBlockSize(maxBlockSize, readHandle,
startEntry, entryBytesWritten);
+
+ try (BlockAwareSegmentInputStream blockStream = new
BlockAwareSegmentInputStreamImpl(
+ readHandle, startEntry, blockSize)) {
+
+ UploadPartResult uploadRes = s3client.uploadPart(
+ new UploadPartRequest()
+ .withBucketName(bucket)
+ .withKey(dataBlockKey)
+ .withUploadId(dataBlockRes.getUploadId())
+ .withInputStream(blockStream)
+ .withPartSize(blockSize)
+ .withPartNumber(partId));
+ etags.add(uploadRes.getPartETag());
+ indexBuilder.addBlock(startEntry, partId, blockSize);
+
+ if (blockStream.getEndEntryId() != -1) {
+ startEntry = blockStream.getEndEntryId() + 1;
+ } else {
+ // could not read entry from ledger.
+ break;
+ }
+ entryBytesWritten +=
blockStream.getBlockEntryBytesCount();
+ partId++;
+ }
}
- });
+
+ s3client.completeMultipartUpload(new
CompleteMultipartUploadRequest()
+ .withBucketName(bucket).withKey(dataBlockKey)
+ .withUploadId(dataBlockRes.getUploadId())
+ .withPartETags(etags));
+ } catch (Throwable t) {
+ s3client.abortMultipartUpload(
+ new AbortMultipartUploadRequest(bucket, dataBlockKey,
dataBlockRes.getUploadId()));
+ promise.completeExceptionally(t);
+ return;
+ }
+
+ // upload index block
+ try (OffloadIndexBlock index = indexBuilder.build();
+ InputStream indexStream = index.toStream()) {
+ // write the index block
+ ObjectMetadata metadata = new ObjectMetadata();
+ metadata.setContentLength(indexStream.available());
+ s3client.putObject(new PutObjectRequest(
+ bucket,
+ indexBlockOffloadKey(readHandle, uuid),
+ indexStream,
+ metadata));
+ promise.complete(null);
+ } catch (Throwable t) {
+ s3client.deleteObject(bucket, dataBlockOffloadKey(readHandle,
uuid));
+ promise.completeExceptionally(t);
+ return;
+ }
+ });
return promise;
}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/BlockAwareSegmentInputStreamImpl.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/BlockAwareSegmentInputStreamImpl.java
index 624b9d96b3..0fdf48f7f8 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/BlockAwareSegmentInputStreamImpl.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/BlockAwareSegmentInputStreamImpl.java
@@ -21,7 +21,6 @@
import static com.google.common.base.Preconditions.checkState;
import com.google.common.collect.Lists;
-import com.google.common.primitives.Ints;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.CompositeByteBuf;
import io.netty.buffer.PooledByteBufAllocator;
@@ -64,7 +63,7 @@
// how many entries want to read from ReadHandle each time.
private static final int ENTRIES_PER_READ = 100;
// buf the entry size and entry id.
- static final int ENTRY_HEADER_SIZE = 4 /* entry size*/ + 8 /* entry id */;
+ static final int ENTRY_HEADER_SIZE = 4 /* entry size */ + 8 /* entry id */;
// Keep a list of all entries ByteBuf, each ByteBuf contains 2 buf: entry
header and entry content.
private List<ByteBuf> entriesByteBuf = null;
@@ -204,5 +203,16 @@ public long getEndEntryId() {
public int getBlockEntryBytesCount() {
return dataBlockFullOffset - DataBlockHeaderImpl.getDataStartOffset()
- ENTRY_HEADER_SIZE * blockEntryCount;
}
+
+ // Calculate the block size after uploaded `entryBytesAlreadyWritten` bytes
+ public static int calculateBlockSize(int maxBlockSize, ReadHandle
readHandle,
+ long firstEntryToWrite, long
entryBytesAlreadyWritten) {
+ return (int)Math.min(
+ maxBlockSize,
+ (readHandle.getLastAddConfirmed() - firstEntryToWrite + 1) *
ENTRY_HEADER_SIZE
+ + (readHandle.getLength() - entryBytesAlreadyWritten)
+ + DataBlockHeaderImpl.getDataStartOffset());
+ }
+
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloaderTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloaderTest.java
index 059db65d27..583b0e03bf 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloaderTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloaderTest.java
@@ -18,9 +18,16 @@
*/
package org.apache.pulsar.broker.s3offload;
+import static
org.apache.pulsar.broker.s3offload.S3ManagedLedgerOffloader.dataBlockOffloadKey;
+import static
org.apache.pulsar.broker.s3offload.S3ManagedLedgerOffloader.indexBlockOffloadKey;
+import static org.mockito.Matchers.any;
+import com.amazonaws.AmazonServiceException;
+import com.amazonaws.services.s3.AmazonS3;
+import com.amazonaws.services.s3.model.S3Object;
import io.netty.util.concurrent.DefaultThreadFactory;
+import java.io.DataInputStream;
import java.util.HashMap;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
@@ -28,20 +35,23 @@
import java.util.concurrent.ScheduledExecutorService;
import org.apache.bookkeeper.client.BookKeeper;
-import org.apache.bookkeeper.client.LedgerHandle;
+import org.apache.bookkeeper.client.BookKeeper.DigestType;
import org.apache.bookkeeper.client.MockBookKeeper;
-import org.apache.bookkeeper.client.api.DigestType;
+import org.apache.bookkeeper.client.MockLedgerHandle;
import org.apache.bookkeeper.client.api.ReadHandle;
import org.apache.bookkeeper.mledger.LedgerOffloader;
-
import org.apache.pulsar.broker.PulsarServerException;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
+import
org.apache.pulsar.broker.s3offload.impl.BlockAwareSegmentInputStreamImpl;
+import org.apache.pulsar.broker.s3offload.impl.DataBlockHeaderImpl;
+import org.apache.pulsar.broker.s3offload.impl.OffloadIndexBlockImpl;
+import
org.apache.pulsar.broker.s3offload.impl.OffloadIndexTest.LedgerMetadataMock;
+import org.mockito.Mockito;
import org.testng.Assert;
import org.testng.annotations.Test;
class S3ManagedLedgerOffloaderTest extends S3TestBase {
-
final ScheduledExecutorService scheduler;
final MockBookKeeper bk;
@@ -50,21 +60,86 @@
bk = new
MockBookKeeper(MockedPulsarServiceBaseTest.createMockZooKeeper());
}
- private ReadHandle buildReadHandle() throws Exception {
- LedgerHandle lh = bk.createLedger(1,1,1, BookKeeper.DigestType.CRC32,
"foobar".getBytes());
- lh.addEntry("foobar".getBytes());
+ private ReadHandle buildReadHandle(int entryCount) throws Exception {
+ MockLedgerHandle lh = (MockLedgerHandle)bk.createLedger(1,1,1,
BookKeeper.DigestType.CRC32, "foobar".getBytes());
+
+ for (int index = 0; index < entryCount; index ++) {
+ lh.addEntry(("foooobarrr").getBytes()); // add entry with 10 bytes
data
+ }
+
lh.close();
- ReadHandle readHandle = bk.newOpenLedgerOp().withLedgerId(lh.getId())
-
.withPassword("foobar".getBytes()).withDigestType(DigestType.CRC32).execute().get();
- return lh;
+ // mock ledgerMetadata, so the lac in metadata is not -1;
+ MockLedgerHandle spy = Mockito.spy(lh);
+ LedgerMetadataMock metadata = new LedgerMetadataMock(1, 1, 1,
+ DigestType.CRC32C, "foobar".getBytes(), null, false);
+ metadata.setLastEntryId(entryCount - 1);
+ Mockito.when(spy.getLedgerMetadata()).thenReturn(metadata);
+
+ return spy;
+ }
+
+ private void verifyS3ObjectRead(S3Object object, S3Object indexObject,
ReadHandle readHandle, int indexEntryCount, int entryCount, int maxBlockSize)
throws Exception {
+ DataInputStream dis = new DataInputStream(object.getObjectContent());
+ int isLength = dis.available();
+
+ // read out index block
+ DataInputStream indexBlockIs = new
DataInputStream(indexObject.getObjectContent());
+ OffloadIndexBlock indexBlock = OffloadIndexBlockImpl.get(indexBlockIs);
+
+ // 1. verify index block with passed in index entry count
+ Assert.assertEquals(indexBlock.getEntryCount(), indexEntryCount);
+
+ // 2. verify index block read out each indexEntry.
+ int entryIdTracker = 0;
+ int startPartIdTracker = 1;
+ int startOffsetTracker = 0;
+ long entryBytesUploaded = 0;
+ int entryLength = 10;
+ for (int i = 0; i < indexEntryCount; i ++) {
+ // 2.1 verify each indexEntry in header block
+ OffloadIndexEntry indexEntry =
indexBlock.getIndexEntryForEntry(entryIdTracker);
+
+ Assert.assertEquals(indexEntry.getPartId(), startPartIdTracker);
+ Assert.assertEquals(indexEntry.getEntryId(), entryIdTracker);
+ Assert.assertEquals(indexEntry.getOffset(), startOffsetTracker);
+
+ // read out and verify each data block related to this index entry
+ // 2.2 verify data block header.
+ DataBlockHeader headerReadout =
DataBlockHeaderImpl.fromStream(dis);
+ int expectedBlockSize = BlockAwareSegmentInputStreamImpl
+ .calculateBlockSize(maxBlockSize, readHandle, entryIdTracker,
entryBytesUploaded);
+ Assert.assertEquals(headerReadout.getBlockLength(),
expectedBlockSize);
+ Assert.assertEquals(headerReadout.getFirstEntryId(),
entryIdTracker);
+
+ // 2.3 verify data block
+ int entrySize = 0;
+ long entryId = 0;
+ for (int bytesReadout = headerReadout.getBlockLength() -
DataBlockHeaderImpl.getDataStartOffset();
+ bytesReadout > 0;
+ bytesReadout -= (4 + 8 + entrySize)) {
+ entrySize = dis.readInt();
+ entryId = dis.readLong();
+ byte[] bytes = new byte[(int) entrySize];
+ dis.read(bytes);
+
+ Assert.assertEquals(entrySize, entryLength);
+ Assert.assertEquals(entryId, entryIdTracker ++);
+ entryBytesUploaded += entrySize;
+ }
+
+ startPartIdTracker ++;
+ startOffsetTracker += headerReadout.getBlockLength();
+ }
+
+ return;
}
@Test
public void testHappyCase() throws Exception {
- LedgerOffloader offloader = new S3ManagedLedgerOffloader(s3client,
BUCKET, scheduler);
+ LedgerOffloader offloader = new S3ManagedLedgerOffloader(s3client,
BUCKET, scheduler, 1024);
- offloader.offload(buildReadHandle(), UUID.randomUUID(), new
HashMap<>()).get();
+ offloader.offload(buildReadHandle(1), UUID.randomUUID(), new
HashMap<>()).get();
}
@Test
@@ -77,7 +152,7 @@ public void testBucketDoesNotExist() throws Exception {
LedgerOffloader offloader = S3ManagedLedgerOffloader.create(conf,
scheduler);
try {
- offloader.offload(buildReadHandle(), UUID.randomUUID(), new
HashMap<>()).get();
+ offloader.offload(buildReadHandle(1), UUID.randomUUID(), new
HashMap<>()).get();
Assert.fail("Shouldn't be able to add to bucket");
} catch (ExecutionException e) {
Assert.assertTrue(e.getMessage().contains("NoSuchBucket"));
@@ -111,5 +186,141 @@ public void testNoBucketConfigured() throws Exception {
// correct
}
}
+
+ @Test
+ public void testOffload() throws Exception {
+ int entryLength = 10;
+ int entryNumberEachBlock = 10;
+ ServiceConfiguration conf = new ServiceConfiguration();
+
conf.setManagedLedgerOffloadDriver(S3ManagedLedgerOffloader.DRIVER_NAME);
+
+ conf.setS3ManagedLedgerOffloadBucket(BUCKET);
+ conf.setS3ManagedLedgerOffloadRegion("eu-west-1");
+ conf.setS3ManagedLedgerOffloadServiceEndpoint(s3endpoint);
+ conf.setS3ManagedLedgerOffloadMaxBlockSizeInBytes(
+ DataBlockHeaderImpl.getDataStartOffset() + (entryLength + 12) *
entryNumberEachBlock);
+ LedgerOffloader offloader = S3ManagedLedgerOffloader.create(conf,
scheduler);
+
+ // offload 30 entries, which will be placed into 3 data blocks.
+ int entryCount = 30;
+ ReadHandle readHandle = buildReadHandle(entryCount);
+ UUID uuid = UUID.randomUUID();
+ offloader.offload(readHandle, uuid, new HashMap<>()).get();
+
+ S3Object obj = s3client.getObject(BUCKET,
dataBlockOffloadKey(readHandle, uuid));
+ S3Object indexObj = s3client.getObject(BUCKET,
S3ManagedLedgerOffloader.indexBlockOffloadKey(readHandle, uuid));
+
+ verifyS3ObjectRead(obj, indexObj, readHandle, 3, 30,
conf.getS3ManagedLedgerOffloadMaxBlockSizeInBytes());
+ }
+
+ @Test
+ public void testOffloadFailInitDataBlockUpload() throws Exception {
+ int maxBlockSize = 1024;
+ int entryCount = 3;
+ ReadHandle readHandle = buildReadHandle(entryCount);
+ UUID uuid = UUID.randomUUID();
+ String failureString = "fail InitDataBlockUpload";
+
+ // mock throw exception when initiateMultipartUpload
+ try {
+ AmazonS3 mockS3client = Mockito.spy(s3client);
+ Mockito
+ .doThrow(new AmazonServiceException(failureString))
+ .when(mockS3client).initiateMultipartUpload(any());
+
+ LedgerOffloader offloader = new
S3ManagedLedgerOffloader(mockS3client, BUCKET, scheduler, maxBlockSize);
+ offloader.offload(readHandle, uuid, new HashMap<>()).get();
+ Assert.fail("Should throw exception when initiateMultipartUpload");
+ } catch (Exception e) {
+ // excepted
+ Assert.assertTrue(e.getCause() instanceof AmazonServiceException);
+
Assert.assertTrue(e.getCause().getMessage().contains(failureString));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
dataBlockOffloadKey(readHandle, uuid)));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
indexBlockOffloadKey(readHandle, uuid)));
+ }
+ }
+
+ @Test
+ public void testOffloadFailDataBlockPartUpload() throws Exception {
+ int maxBlockSize = 1024;
+ int entryCount = 3;
+ ReadHandle readHandle = buildReadHandle(entryCount);
+ UUID uuid = UUID.randomUUID();
+ String failureString = "fail DataBlockPartUpload";
+
+ // mock throw exception when uploadPart
+ try {
+ AmazonS3 mockS3client = Mockito.spy(s3client);
+ Mockito
+ .doThrow(new AmazonServiceException("fail
DataBlockPartUpload"))
+ .when(mockS3client).uploadPart(any());
+ Mockito.doNothing().when(mockS3client).abortMultipartUpload(any());
+
+ LedgerOffloader offloader = new
S3ManagedLedgerOffloader(mockS3client, BUCKET, scheduler, maxBlockSize);
+ offloader.offload(readHandle, uuid, new HashMap<>()).get();
+ Assert.fail("Should throw exception for when uploadPart");
+ } catch (Exception e) {
+ // excepted
+ Assert.assertTrue(e.getCause() instanceof AmazonServiceException);
+
Assert.assertTrue(e.getCause().getMessage().contains(failureString));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
dataBlockOffloadKey(readHandle, uuid)));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
indexBlockOffloadKey(readHandle, uuid)));
+ }
+ }
+
+ @Test
+ public void testOffloadFailDataBlockUploadComplete() throws Exception {
+ int maxBlockSize = 1024;
+ int entryCount = 3;
+ ReadHandle readHandle = buildReadHandle(entryCount);
+ UUID uuid = UUID.randomUUID();
+ String failureString = "fail DataBlockUploadComplete";
+
+ // mock throw exception when completeMultipartUpload
+ try {
+ AmazonS3 mockS3client = Mockito.spy(s3client);
+ Mockito
+ .doThrow(new AmazonServiceException(failureString))
+ .when(mockS3client).completeMultipartUpload(any());
+ Mockito.doNothing().when(mockS3client).abortMultipartUpload(any());
+
+ LedgerOffloader offloader = new
S3ManagedLedgerOffloader(mockS3client, BUCKET, scheduler, maxBlockSize);
+ offloader.offload(readHandle, uuid, new HashMap<>()).get();
+ Assert.fail("Should throw exception for when
completeMultipartUpload");
+ } catch (Exception e) {
+ // excepted
+ Assert.assertTrue(e.getCause() instanceof AmazonServiceException);
+
Assert.assertTrue(e.getCause().getMessage().contains(failureString));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
dataBlockOffloadKey(readHandle, uuid)));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
indexBlockOffloadKey(readHandle, uuid)));
+ }
+ }
+
+ @Test
+ public void testOffloadFailPutIndexBlock() throws Exception {
+ int maxBlockSize = 1024;
+ int entryCount = 3;
+ ReadHandle readHandle = buildReadHandle(entryCount);
+ UUID uuid = UUID.randomUUID();
+ String failureString = "fail putObject";
+
+ // mock throw exception when putObject
+ try {
+ AmazonS3 mockS3client = Mockito.spy(s3client);
+ Mockito
+ .doThrow(new AmazonServiceException(failureString))
+ .when(mockS3client).putObject(any());
+
+ LedgerOffloader offloader = new
S3ManagedLedgerOffloader(mockS3client, BUCKET, scheduler, maxBlockSize);
+ offloader.offload(readHandle, uuid, new HashMap<>()).get();
+ Assert.fail("Should throw exception for when putObject for index
block");
+ } catch (Exception e) {
+ // excepted
+ Assert.assertTrue(e.getCause() instanceof AmazonServiceException);
+
Assert.assertTrue(e.getCause().getMessage().contains(failureString));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
dataBlockOffloadKey(readHandle, uuid)));
+ Assert.assertFalse(s3client.doesObjectExist(BUCKET,
indexBlockOffloadKey(readHandle, uuid)));
+ }
+ }
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexTest.java
index 3867899b57..916c4cda57 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexTest.java
@@ -62,7 +62,7 @@ public void offloadIndexEntryImplTest() {
// use mock to setLastEntryId
- class LedgerMetadataMock extends
org.apache.bookkeeper.client.LedgerMetadata {
+ public static class LedgerMetadataMock extends
org.apache.bookkeeper.client.LedgerMetadata {
long lastId = 0;
public LedgerMetadataMock(int ensembleSize, int writeQuorumSize, int
ackQuorumSize, org.apache.bookkeeper.client.BookKeeper.DigestType digestType,
byte[] password, Map<String, byte[]> customMetadata, boolean
storeSystemtimeAsLedgerCreationTime) {
super(ensembleSize, writeQuorumSize, ackQuorumSize, digestType,
password, customMetadata, storeSystemtimeAsLedgerCreationTime);
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services