This is an automated email from the ASF dual-hosted git repository.
sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new c4edead PIP-17: impl offload() for S3ManagedLedgerOffloader (#1746)
c4edead is described below
commit c4edead69220388bfb855c8d32442745eddd03b9
Author: Jia Zhai <[email protected]>
AuthorDate: Tue May 15 19:58:00 2018 +0800
PIP-17: impl offload() for S3ManagedLedgerOffloader (#1746)
* add write for S3ManagedLedgerOffloader
* merge master, change following comments
* change following @ivan's comments
* change following @ivan's comments
* fix conf error
---
conf/broker.conf | 3 +
.../apache/pulsar/broker/ServiceConfiguration.java | 12 ++
.../broker/s3offload/S3ManagedLedgerOffloader.java | 128 +++++++++--
.../impl/BlockAwareSegmentInputStreamImpl.java | 14 +-
.../s3offload/S3ManagedLedgerOffloaderTest.java | 237 +++++++++++++++++++--
.../broker/s3offload/impl/OffloadIndexTest.java | 2 +-
6 files changed, 367 insertions(+), 29 deletions(-)
diff --git a/conf/broker.conf b/conf/broker.conf
index e431810..cf9faa8 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 43d929d..74363cc 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 @@ public class ServiceConfiguration implements
PulsarConfiguration {
// 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 class ServiceConfiguration implements
PulsarConfiguration {
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 163b79b..76920de 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 class S3ManagedLedgerOffloader implements
LedgerOffloader {
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 class S3ManagedLedgerOffloader implements
LedgerOffloader {
} 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 624b9d9..0fdf48f 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 @@ package org.apache.pulsar.broker.s3offload.impl;
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 @@ public class BlockAwareSegmentInputStreamImpl extends
BlockAwareSegmentInputStre
// 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 class BlockAwareSegmentInputStreamImpl extends
BlockAwareSegmentInputStre
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 059db65..583b0e0 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.Executors;
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 @@ class S3ManagedLedgerOffloaderTest extends S3TestBase {
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 @@ class S3ManagedLedgerOffloaderTest extends S3TestBase {
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 @@ class S3ManagedLedgerOffloaderTest extends S3TestBase {
// 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 3867899..916c4cd 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 class OffloadIndexTest {
// 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);
--
To stop receiving notification emails like this one, please contact
[email protected].