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].

Reply via email to