This is an automated email from the ASF dual-hosted git repository.
clintropolis pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 15051a5c09b feat: add support for push/pull/kill segments in google
cloud deep storage without zip compression (#20382)
15051a5c09b is described below
commit 15051a5c09b04eda967dee9dbc552a6dc658d794
Author: Clint Wylie <[email protected]>
AuthorDate: Tue Sep 22 11:25:03 2026 -0700
feat: add support for push/pull/kill segments in google cloud deep storage
without zip compression (#20382)
---
docs/development/extensions-core/google.md | 1 +
.../storage/google/GoogleDataSegmentKiller.java | 39 ++++++-
.../storage/google/GoogleDataSegmentPuller.java | 60 +++++++++--
.../storage/google/GoogleDataSegmentPusher.java | 112 ++++++++++++++++++++-
.../google/GoogleTimestampVersionedDataFinder.java | 4 +-
.../google/GoogleDataSegmentKillerTest.java | 32 ++++++
.../google/GoogleDataSegmentPullerTest.java | 52 +++++++++-
.../google/GoogleDataSegmentPusherTest.java | 108 +++++++++++++++++++-
.../GoogleTimestampVersionedDataFinderTest.java | 4 +-
9 files changed, 385 insertions(+), 27 deletions(-)
diff --git a/docs/development/extensions-core/google.md
b/docs/development/extensions-core/google.md
index bff1d269611..8d95e1b48be 100644
--- a/docs/development/extensions-core/google.md
+++ b/docs/development/extensions-core/google.md
@@ -53,3 +53,4 @@ To configure connectivity to google cloud, run druid
processes with `GOOGLE_APPL
|`druid.google.bucket`||Google Storage bucket name.|Must be set.|
|`druid.google.prefix`|A prefix string that will be prepended to the blob
names for the segments published to Google deep storage| |""|
|`druid.google.maxListingLength`|maximum number of input files matching a
given prefix to retrieve at a time| |1024|
+|`druid.storage.zip`|Whether segments are written as directories of objects
(`false`) or as a single `index.zip` object (`true`). Writing segments unzipped
requires permission to delete objects under the segment path, so that a push
replaces the segment already there.|`true`, `false`|`true`|
diff --git
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentKiller.java
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentKiller.java
index 42cc7c194b4..e4730ced127 100644
---
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentKiller.java
+++
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentKiller.java
@@ -61,19 +61,48 @@ public class GoogleDataSegmentKiller implements
DataSegmentKiller
Map<String, Object> loadSpec = segment.getLoadSpec();
final String bucket = MapUtils.getString(loadSpec, "bucket");
final String indexPath = MapUtils.getString(loadSpec, "path");
- final String descriptorPath = DataSegmentKiller.descriptorPath(indexPath);
try {
- deleteIfPresent(bucket, indexPath);
- // descriptor.json is a file to store segment metadata in deep storage.
This file is deprecated and not stored
- // anymore, but we still delete them if exists.
- deleteIfPresent(bucket, descriptorPath);
+ if (indexPath.endsWith("/")) {
+ // segment was pushed unzipped, so the path names a directory of
objects; delete every one of them
+ deleteObjectsInPath(bucket, indexPath);
+ } else {
+ deleteIfPresent(bucket, indexPath);
+ // descriptor.json is a file to store segment metadata in deep
storage. This file is deprecated and not stored
+ // anymore, but we still delete them if exists.
+ deleteIfPresent(bucket, DataSegmentKiller.descriptorPath(indexPath));
+ }
}
catch (StorageException e) {
throw new SegmentLoadingException(e, "Couldn't kill segment[%s]: [%s]",
segment.getId(), e.getMessage());
}
}
+ private void deleteObjectsInPath(String bucket, String pathPrefix) throws
SegmentLoadingException
+ {
+ try {
+ GoogleUtils.deleteObjectsInPath(
+ storage,
+ inputDataConfig,
+ bucket,
+ pathPrefix,
+ Predicates.alwaysTrue()
+ );
+ }
+ catch (StorageException e) {
+ throw e;
+ }
+ catch (Exception e) {
+ throw new SegmentLoadingException(
+ e,
+ "Couldn't delete objects under [gs://%s/%s]: [%s]",
+ bucket,
+ pathPrefix,
+ e.getMessage()
+ );
+ }
+ }
+
private void deleteIfPresent(String bucket, String path)
{
try {
diff --git
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPuller.java
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPuller.java
index 8e155f6a4ce..f731faaaeb4 100644
---
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPuller.java
+++
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPuller.java
@@ -20,8 +20,11 @@
package org.apache.druid.storage.google;
import com.google.common.base.Predicate;
+import com.google.common.collect.ImmutableList;
import com.google.inject.Inject;
+import org.apache.druid.data.input.impl.CloudObjectLocation;
import org.apache.druid.java.util.common.FileUtils;
+import org.apache.druid.java.util.common.RetryUtils;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.logger.Logger;
import org.apache.druid.segment.loading.SegmentLoadingException;
@@ -32,17 +35,21 @@ import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.net.URI;
+import java.nio.file.Paths;
+import java.util.Iterator;
public class GoogleDataSegmentPuller implements URIDataPuller
{
private static final Logger LOG = new Logger(GoogleDataSegmentPuller.class);
protected final GoogleStorage storage;
+ private final GoogleInputDataConfig inputDataConfig;
@Inject
- public GoogleDataSegmentPuller(final GoogleStorage storage)
+ public GoogleDataSegmentPuller(final GoogleStorage storage, final
GoogleInputDataConfig inputDataConfig)
{
this.storage = storage;
+ this.inputDataConfig = inputDataConfig;
}
FileUtils.FileCopyResult getSegmentFiles(final String bucket, final String
path, File outDir)
@@ -53,13 +60,12 @@ public class GoogleDataSegmentPuller implements
URIDataPuller
try {
FileUtils.mkdirp(outDir);
- final GoogleByteSource byteSource = new GoogleByteSource(storage,
bucket, path);
- final FileUtils.FileCopyResult result = CompressionUtils.unzip(
- byteSource,
- outDir,
- GoogleUtils::isRetryable,
- false
- );
+ // A trailing slash means the segment was pushed unzipped
(druid.storage.zip=false), so the path names a
+ // directory of objects to pull individually rather than a single zip to
unpack.
+ final FileUtils.FileCopyResult result = path.endsWith("/")
+ ?
getSegmentFilesFromDirectory(bucket, path, outDir)
+ : unzipSegmentFiles(bucket,
path, outDir);
+
LOG.info("Loaded %d bytes from [%s] to [%s]", result.size(), path,
outDir.getAbsolutePath());
return result;
}
@@ -79,6 +85,44 @@ public class GoogleDataSegmentPuller implements URIDataPuller
}
}
+ private FileUtils.FileCopyResult unzipSegmentFiles(final String bucket,
final String path, final File outDir)
+ throws IOException
+ {
+ final GoogleByteSource byteSource = new GoogleByteSource(storage, bucket,
path);
+ return CompressionUtils.unzip(
+ byteSource,
+ outDir,
+ GoogleUtils::isRetryable,
+ false
+ );
+ }
+
+ private FileUtils.FileCopyResult getSegmentFilesFromDirectory(
+ final String bucket,
+ final String pathPrefix,
+ final File outDir
+ )
+ {
+ final Iterator<GoogleStorageObjectMetadata> objects =
GoogleUtils.lazyFetchingStorageObjectsIterator(
+ storage,
+ ImmutableList.of(new CloudObjectLocation(bucket,
pathPrefix).toUri(GoogleStorageDruidModule.SCHEME_GS))
+ .iterator(),
+ inputDataConfig.getMaxListingLength()
+ );
+
+ final FileUtils.FileCopyResult copyResult = new FileUtils.FileCopyResult();
+ while (objects.hasNext()) {
+ final GoogleStorageObjectMetadata object = objects.next();
+ final GoogleByteSource byteSource = new GoogleByteSource(storage,
bucket, object.getName());
+ final File outFile = new File(outDir,
Paths.get(object.getName()).getFileName().toString());
+ copyResult.addFiles(
+ FileUtils.retryCopy(byteSource, outFile, GoogleUtils.GOOGLE_RETRY,
RetryUtils.DEFAULT_MAX_TRIES).getFiles()
+ );
+ }
+
+ return copyResult;
+ }
+
@Override
public InputStream getInputStream(URI uri) throws IOException
{
diff --git
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPusher.java
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPusher.java
index 1da95ad1aa5..144e43a1331 100644
---
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPusher.java
+++
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleDataSegmentPusher.java
@@ -24,11 +24,13 @@ import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Strings;
import com.google.common.collect.ImmutableMap;
import com.google.inject.Inject;
+import org.apache.druid.java.util.common.IOE;
import org.apache.druid.java.util.common.RE;
import org.apache.druid.java.util.common.RetryUtils;
import org.apache.druid.java.util.common.logger.Logger;
import org.apache.druid.segment.SegmentUtils;
import org.apache.druid.segment.loading.DataSegmentPusher;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
import org.apache.druid.timeline.DataSegment;
import org.apache.druid.utils.CompressionUtils;
@@ -36,23 +38,39 @@ import java.io.File;
import java.io.IOException;
import java.net.URI;
import java.nio.file.Files;
+import java.util.HashSet;
import java.util.Map;
+import java.util.Set;
public class GoogleDataSegmentPusher implements DataSegmentPusher
{
private static final Logger log = new Logger(GoogleDataSegmentPusher.class);
+ static final String INDEX_ZIP_FILE_NAME = "index.zip";
+
+ /**
+ * Originally Google deep storage always wrote segments zipped, so that is
what {@code druid.storage.zip} falls back
+ * to.
+ */
+ private static final boolean DEFAULT_ZIP = true;
+
private final GoogleStorage storage;
private final GoogleAccountConfig config;
+ private final GoogleInputDataConfig inputDataConfig;
+ private final boolean zip;
@Inject
public GoogleDataSegmentPusher(
final GoogleStorage storage,
- final GoogleAccountConfig config
+ final GoogleAccountConfig config,
+ final GoogleInputDataConfig inputDataConfig,
+ final DeepStorageSegmentConfig deepStorageConfig
)
{
this.storage = storage;
this.config = config;
+ this.inputDataConfig = inputDataConfig;
+ this.zip = deepStorageConfig.isZip(DEFAULT_ZIP);
}
public void insert(final File file, final String contentType, final String
path)
@@ -91,12 +109,29 @@ public class GoogleDataSegmentPusher implements
DataSegmentPusher
public DataSegment pushToPath(File indexFilesDir, DataSegment segment,
String storageDirSuffix) throws IOException
{
final int version = SegmentUtils.getVersionFromDir(indexFilesDir);
+ final String basePath = buildPath(storageDirSuffix);
+
+ try {
+ if (zip) {
+ return pushZip(indexFilesDir, segment, version, basePath);
+ } else {
+ return pushNoZip(indexFilesDir, segment, version, basePath);
+ }
+ }
+ catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private DataSegment pushZip(File indexFilesDir, DataSegment segment, int
version, String basePath)
+ throws IOException
+ {
File indexFile = null;
try {
indexFile = Files.createTempFile("index", ".zip").toFile();
final long indexSize = CompressionUtils.zip(indexFilesDir, indexFile);
- final String indexPath = buildPath(storageDirSuffix + "/" + "index.zip");
+ final String indexPath = basePath + "/" + INDEX_ZIP_FILE_NAME;
final DataSegment outSegment = segment
.withSize(indexSize)
@@ -107,9 +142,6 @@ public class GoogleDataSegmentPusher implements
DataSegmentPusher
return outSegment;
}
- catch (Exception e) {
- throw new RuntimeException(e);
- }
finally {
if (indexFile != null) {
log.debug("Deleting file [%s]", indexFile);
@@ -118,6 +150,76 @@ public class GoogleDataSegmentPusher implements
DataSegmentPusher
}
}
+ /**
+ * Uploads the segment files as they are, one object per file, under {@code
basePath}. The resulting loadSpec path is
+ * the directory itself, with a trailing slash to tell {@link
GoogleDataSegmentPuller} and
+ * {@link GoogleDataSegmentKiller} that it names a directory of files rather
than a single object.
+ */
+ private DataSegment pushNoZip(File indexFilesDir, DataSegment segment, int
version, String basePath)
+ throws IOException
+ {
+ final File[] files = indexFilesDir.listFiles();
+ if (files == null) {
+ throw new IOE("Cannot list directory [%s]", indexFilesDir);
+ }
+
+ final String dirPath = basePath + "/";
+ final Set<String> pushedPaths = new HashSet<>();
+
+ long size = 0;
+ for (final File file : files) {
+ if (file.isFile()) {
+ size += file.length();
+ final String path = dirPath + file.getName();
+ insert(file, "application/octet-stream", path);
+ pushedPaths.add(path);
+ } else {
+ // Segment directories are expected to be flat.
+ throw new IOE("Unexpected subdirectory [%s]", file.getName());
+ }
+ }
+
+ deleteStaleObjects(dirPath, pushedPaths);
+
+ return segment.withSize(size)
+ .withLoadSpec(makeLoadSpec(config.getBucket(), dirPath))
+ .withBinaryVersion(version);
+ }
+
+ /**
+ * Removes everything under {@code dirPath} that this push did not write.
+ * <p>
+ * A zipped push replaces the previous segment outright, because one {@code
index.zip} object overwrites another, and
+ * {@link #push} with {@code useUniquePath = false} is expected to replace a
previous push the same way. Uploading
+ * file by file only overwrites the names the new segment happens to share,
so without this a re-push could leave
+ * behind objects of whatever was there before: a stale {@code index.zip}
from a zipped push, or smoosh chunks from a
+ * larger prior v9 segment.
+ * <p>
+ * Note that this makes an unzipped push require permission to delete
objects under the segment path, which a zipped
+ * push does not.
+ */
+ private void deleteStaleObjects(final String dirPath, final Set<String>
pushedPaths) throws IOException
+ {
+ try {
+ GoogleUtils.deleteObjectsInPath(
+ storage,
+ inputDataConfig,
+ config.getBucket(),
+ dirPath,
+ object -> !pushedPaths.contains(object.getName())
+ );
+ }
+ catch (Exception e) {
+ throw new IOE(
+ e,
+ "Could not remove objects left under [gs://%s/%s] by a previous
push, which would be loaded as part of this"
+ + " segment",
+ config.getBucket(),
+ dirPath
+ );
+ }
+ }
+
@VisibleForTesting
String buildPath(final String path)
{
diff --git
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinder.java
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinder.java
index d95577d5be6..a5dbb59dfc9 100644
---
a/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinder.java
+++
b/extensions-core/google-extensions/src/main/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinder.java
@@ -35,9 +35,9 @@ public class GoogleTimestampVersionedDataFinder extends
GoogleDataSegmentPuller
private static final long MAX_LISTING_KEYS = 1000;
@Inject
- public GoogleTimestampVersionedDataFinder(final GoogleStorage storage)
+ public GoogleTimestampVersionedDataFinder(final GoogleStorage storage, final
GoogleInputDataConfig inputDataConfig)
{
- super(storage);
+ super(storage, inputDataConfig);
}
@Override
diff --git
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentKillerTest.java
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentKillerTest.java
index a83e24fd4eb..41a55af1ac6 100644
---
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentKillerTest.java
+++
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentKillerTest.java
@@ -95,6 +95,38 @@ public class GoogleDataSegmentKillerTest extends
EasyMockSupport
verifyAll();
}
+ @Test
+ public void test_kill_unzippedSegment_deletesEveryObjectInTheDirectory()
throws SegmentLoadingException, IOException
+ {
+ // pushed with druid.storage.zip=false, so the path is the directory
holding the segment files
+ final String unzippedPath =
"test/2015-04-12T00:00:00.000Z_2015-04-13T00:00:00.000Z/1/0/";
+ final DataSegment segment = DataSegment.builder(DATA_SEGMENT)
+ .loadSpec(ImmutableMap.of("bucket",
BUCKET, "path", unzippedPath))
+ .build();
+
+
EasyMock.expect(inputDataConfig.getMaxListingLength()).andReturn(MAX_KEYS).anyTimes();
+ EasyMock.expect(storage.list(EasyMock.eq(BUCKET),
EasyMock.eq(unzippedPath), EasyMock.anyObject(), EasyMock.anyObject()))
+ .andReturn(new GoogleStorageObjectPage(
+ ImmutableList.of(
+ new GoogleStorageObjectMetadata(BUCKET, unzippedPath +
"version.bin", 4L, TIME_0),
+ new GoogleStorageObjectMetadata(BUCKET, unzippedPath +
"meta.smoosh", 8L, TIME_0)
+ ),
+ null
+ ));
+ storage.delete(BUCKET, unzippedPath + "version.bin");
+ EasyMock.expectLastCall();
+ storage.delete(BUCKET, unzippedPath + "meta.smoosh");
+ EasyMock.expectLastCall();
+
+ replayAll();
+
+ GoogleDataSegmentKiller killer = new GoogleDataSegmentKiller(storage,
accountConfig, inputDataConfig);
+
+ killer.kill(segment);
+
+ verifyAll();
+ }
+
@Test
public void killWithErrorTest()
{
diff --git
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPullerTest.java
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPullerTest.java
index 0e7baf02855..1c753c8402f 100644
---
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPullerTest.java
+++
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPullerTest.java
@@ -22,6 +22,7 @@ package org.apache.druid.storage.google;
import com.google.api.client.googleapis.json.GoogleJsonResponseException;
import
com.google.api.client.googleapis.testing.json.GoogleJsonResponseExceptionFactoryTesting;
import com.google.api.client.json.jackson2.JacksonFactory;
+import com.google.common.collect.ImmutableList;
import org.apache.druid.java.util.common.FileUtils;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.segment.loading.SegmentLoadingException;
@@ -30,6 +31,7 @@ import org.easymock.EasyMockSupport;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import java.io.ByteArrayInputStream;
import java.io.File;
import java.io.IOException;
import java.io.InputStream;
@@ -41,6 +43,9 @@ public class GoogleDataSegmentPullerTest extends
EasyMockSupport
{
private static final String BUCKET = "bucket";
private static final String PATH = "/path/to/storage/index.zip";
+ // no leading slash: GoogleDataSegmentPusher.makeLoadSpec strips it, and the
listing normalizes it away anyway
+ private static final String UNZIPPED_PATH = "path/to/storage/";
+ private static final GoogleInputDataConfig INPUT_DATA_CONFIG = new
GoogleInputDataConfig();
@Test
public void testDeleteOutputDirectoryWhenErrorIsRaisedPullingSegmentFiles()
@@ -58,7 +63,7 @@ public class GoogleDataSegmentPullerTest extends
EasyMockSupport
replayAll();
- GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage);
+ GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage,
INPUT_DATA_CONFIG);
puller.getSegmentFiles(BUCKET, PATH, outDir);
Assertions.assertFalse(outDir.exists());
@@ -71,6 +76,47 @@ public class GoogleDataSegmentPullerTest extends
EasyMockSupport
});
}
+ @Test
+ public void testGetSegmentFilesUnzippedSegmentPullsEachObject() throws
IOException, SegmentLoadingException
+ {
+ final File outDir = FileUtils.createTempDir();
+ try {
+ final String versionObject = UNZIPPED_PATH + "version.bin";
+ final String smooshObject = UNZIPPED_PATH + "meta.smoosh";
+ final byte[] versionData = new byte[]{0x0, 0x0, 0x0, 0x9};
+ final byte[] smooshData = new byte[]{0x1, 0x2};
+
+ GoogleStorage storage = createMock(GoogleStorage.class);
+ EasyMock.expect(storage.list(EasyMock.eq(BUCKET),
EasyMock.eq(UNZIPPED_PATH), EasyMock.anyObject(), EasyMock.anyObject()))
+ .andReturn(new GoogleStorageObjectPage(
+ ImmutableList.of(
+ new GoogleStorageObjectMetadata(BUCKET, versionObject,
(long) versionData.length, 0L),
+ new GoogleStorageObjectMetadata(BUCKET, smooshObject,
(long) smooshData.length, 0L)
+ ),
+ null
+ ));
+ EasyMock.expect(storage.getInputStream(BUCKET, versionObject))
+ .andReturn(new ByteArrayInputStream(versionData));
+ EasyMock.expect(storage.getInputStream(BUCKET, smooshObject))
+ .andReturn(new ByteArrayInputStream(smooshData));
+
+ replayAll();
+
+ GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage,
INPUT_DATA_CONFIG);
+ FileUtils.FileCopyResult result = puller.getSegmentFiles(BUCKET,
UNZIPPED_PATH, outDir);
+
+ // the objects land in outDir under their own names, not unpacked from a
zip
+ Assertions.assertTrue(new File(outDir, "version.bin").exists());
+ Assertions.assertTrue(new File(outDir, "meta.smoosh").exists());
+ Assertions.assertEquals(versionData.length + smooshData.length,
result.size());
+
+ verifyAll();
+ }
+ finally {
+ FileUtils.deleteDirectory(outDir);
+ }
+ }
+
@Test
public void testGetVersionBucketNameWithUnderscores() throws IOException
{
@@ -82,7 +128,7 @@ public class GoogleDataSegmentPullerTest extends
EasyMockSupport
EasyMock.expect(storage.version(EasyMock.eq(bucket),
EasyMock.eq(prefix))).andReturn("0");
EasyMock.replay(storage);
- GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage);
+ GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage,
INPUT_DATA_CONFIG);
String actual =
puller.getVersion(URI.create(StringUtils.format("gs://%s/%s", bucket, prefix)));
Assertions.assertEquals(version, actual);
@@ -99,7 +145,7 @@ public class GoogleDataSegmentPullerTest extends
EasyMockSupport
EasyMock.expect(storage.getInputStream(EasyMock.eq(bucket),
EasyMock.eq(prefix))).andReturn(EasyMock.createMock(InputStream.class));
EasyMock.replay(storage);
- GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage);
+ GoogleDataSegmentPuller puller = new GoogleDataSegmentPuller(storage,
INPUT_DATA_CONFIG);
puller.getInputStream(URI.create(StringUtils.format("gs://%s/%s", bucket,
prefix)));
EasyMock.verify(storage);
diff --git
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPusherTest.java
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPusherTest.java
index b74b77d8239..012c4c3c1c8 100644
---
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPusherTest.java
+++
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleDataSegmentPusherTest.java
@@ -19,11 +19,15 @@
package org.apache.druid.storage.google;
+import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.io.Files;
import org.apache.druid.java.util.common.Intervals;
+import org.apache.druid.segment.loading.DeepStorageSegmentConfig;
import org.apache.druid.timeline.DataSegment;
+import org.apache.druid.timeline.SegmentId;
import org.apache.druid.timeline.partition.NoneShardSpec;
+import org.apache.druid.timeline.partition.NumberedShardSpec;
import org.easymock.EasyMock;
import org.easymock.EasyMockSupport;
import org.junit.jupiter.api.Assertions;
@@ -42,6 +46,9 @@ public class GoogleDataSegmentPusherTest extends
EasyMockSupport
private static final String BUCKET = "bucket";
private static final String PREFIX = "prefix";
+ private static final GoogleInputDataConfig INPUT_DATA_CONFIG = new
GoogleInputDataConfig();
+ private static final DeepStorageSegmentConfig ZIP_CONFIG = new
DeepStorageSegmentConfig(true);
+ private static final DeepStorageSegmentConfig NO_ZIP_CONFIG = new
DeepStorageSegmentConfig(false);
private GoogleStorage storage;
private GoogleAccountConfig googleAccountConfig;
@@ -78,7 +85,7 @@ public class GoogleDataSegmentPusherTest extends
EasyMockSupport
);
GoogleDataSegmentPusher pusher =
createMockBuilder(GoogleDataSegmentPusher.class)
- .withConstructor(storage, googleAccountConfig)
+ .withConstructor(storage, googleAccountConfig, INPUT_DATA_CONFIG,
ZIP_CONFIG)
.addMockedMethod("insert", File.class, String.class, String.class)
.createMock();
@@ -107,6 +114,95 @@ public class GoogleDataSegmentPusherTest extends
EasyMockSupport
verifyAll();
}
+ @Test
+ public void testPushNoZip() throws Exception
+ {
+ final byte[] data = new byte[]{0x0, 0x0, 0x0, 0x1};
+ Files.write(data, new File(tempFolder, "version.bin"));
+ Files.write(data, new File(tempFolder, "meta.smoosh"));
+
+ DataSegment segmentToPush = newSegmentToPush(2 * data.length);
+
+ GoogleDataSegmentPusher pusher =
createMockBuilder(GoogleDataSegmentPusher.class)
+ .withConstructor(storage, googleAccountConfig, INPUT_DATA_CONFIG,
NO_ZIP_CONFIG)
+ .addMockedMethod("insert", File.class, String.class, String.class)
+ .createMock();
+
+ final String expectedDir = PREFIX + "/" +
pusher.getStorageDir(segmentToPush, false);
+ pusher.insert(EasyMock.anyObject(File.class), EasyMock.anyString(),
EasyMock.eq(expectedDir + "/version.bin"));
+ EasyMock.expectLastCall();
+ pusher.insert(EasyMock.anyObject(File.class), EasyMock.anyString(),
EasyMock.eq(expectedDir + "/meta.smoosh"));
+ EasyMock.expectLastCall();
+
+ // nothing else is under the path, so there is nothing to clean up
+ EasyMock.expect(storage.list(EasyMock.eq(BUCKET), EasyMock.eq(expectedDir
+ "/"), EasyMock.anyObject(), EasyMock.anyObject()))
+ .andReturn(new GoogleStorageObjectPage(
+ ImmutableList.of(
+ new GoogleStorageObjectMetadata(BUCKET, expectedDir +
"/version.bin", (long) data.length, 0L),
+ new GoogleStorageObjectMetadata(BUCKET, expectedDir +
"/meta.smoosh", (long) data.length, 0L)
+ ),
+ null
+ ));
+
+ replayAll();
+
+ DataSegment segment = pusher.push(tempFolder, segmentToPush, false);
+
+ // the trailing slash is what marks the path as a directory of files
rather than a single object
+ Assertions.assertEquals(ImmutableMap.of(
+ "type", GoogleStorageDruidModule.SCHEME,
+ "bucket", BUCKET,
+ "path", expectedDir + "/"
+ ), segment.getLoadSpec());
+ Assertions.assertEquals(2 * data.length, segment.getSize());
+
+ verifyAll();
+ }
+
+ @Test
+ public void testPushNoZipRemovesObjectsLeftByPreviousPush() throws Exception
+ {
+ final byte[] data = new byte[]{0x0, 0x0, 0x0, 0x1};
+ Files.write(data, new File(tempFolder, "version.bin"));
+
+ DataSegment segmentToPush = newSegmentToPush(data.length);
+
+ GoogleDataSegmentPusher pusher =
createMockBuilder(GoogleDataSegmentPusher.class)
+ .withConstructor(storage, googleAccountConfig, INPUT_DATA_CONFIG,
NO_ZIP_CONFIG)
+ .addMockedMethod("insert", File.class, String.class, String.class)
+ .createMock();
+
+ final String expectedDir = PREFIX + "/" +
pusher.getStorageDir(segmentToPush, false);
+ pusher.insert(EasyMock.anyObject(File.class), EasyMock.anyString(),
EasyMock.eq(expectedDir + "/version.bin"));
+ EasyMock.expectLastCall();
+
+ // an index.zip from a zipped push and a smoosh chunk from a larger prior
segment: nothing in the new segment
+ // references either, but the puller would list and download both
+ final String staleZip = expectedDir + "/index.zip";
+ final String staleChunk = expectedDir + "/00001.smoosh";
+ EasyMock.expect(storage.list(EasyMock.eq(BUCKET), EasyMock.eq(expectedDir
+ "/"), EasyMock.anyObject(), EasyMock.anyObject()))
+ .andReturn(new GoogleStorageObjectPage(
+ ImmutableList.of(
+ new GoogleStorageObjectMetadata(BUCKET, expectedDir +
"/version.bin", (long) data.length, 0L),
+ new GoogleStorageObjectMetadata(BUCKET, staleZip, 128L,
0L),
+ new GoogleStorageObjectMetadata(BUCKET, staleChunk, 64L,
0L)
+ ),
+ null
+ ));
+
+ // only the objects this push did not write are removed
+ storage.delete(BUCKET, staleZip);
+ EasyMock.expectLastCall();
+ storage.delete(BUCKET, staleChunk);
+ EasyMock.expectLastCall();
+
+ replayAll();
+
+ pusher.push(tempFolder, segmentToPush, false);
+
+ verifyAll();
+ }
+
@Test
public void testBuildPath()
{
@@ -114,10 +210,18 @@ public class GoogleDataSegmentPusherTest extends
EasyMockSupport
StringBuilder sb = new StringBuilder();
sb.setLength(0);
config.setPrefix(sb.toString()); // avoid cached empty string
- GoogleDataSegmentPusher pusher = new GoogleDataSegmentPusher(storage,
config);
+ GoogleDataSegmentPusher pusher = new GoogleDataSegmentPusher(storage,
config, INPUT_DATA_CONFIG, ZIP_CONFIG);
Assertions.assertEquals("/path", pusher.buildPath("/path"));
config.setPrefix(null);
Assertions.assertEquals("/path", pusher.buildPath("/path"));
}
+ private static DataSegment newSegmentToPush(long size)
+ {
+ return DataSegment.builder(SegmentId.of("foo", Intervals.of("2015/2016"),
"0", 0))
+ .shardSpec(new NumberedShardSpec(0, 1))
+ .binaryVersion(0)
+ .size(size)
+ .build();
+ }
}
diff --git
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinderTest.java
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinderTest.java
index ecf79b01806..02af5191a42 100644
---
a/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinderTest.java
+++
b/extensions-core/google-extensions/src/test/java/org/apache/druid/storage/google/GoogleTimestampVersionedDataFinderTest.java
@@ -46,7 +46,7 @@ public class GoogleTimestampVersionedDataFinderTest
storageObject4.setLastUpdateTimeMillis(System.currentTimeMillis() + 100);
final GoogleStorage storage =
ObjectStorageIteratorTest.makeMockClient(ImmutableList.of(storageObject1,
storageObject2, storageObject3, storageObject4));
- final GoogleTimestampVersionedDataFinder finder = new
GoogleTimestampVersionedDataFinder(storage);
+ final GoogleTimestampVersionedDataFinder finder = new
GoogleTimestampVersionedDataFinder(storage, new GoogleInputDataConfig());
Pattern pattern = Pattern.compile("v.*");
URI latest =
finder.getLatestVersion(URI.create(StringUtils.format("gs://%s/%s", bucket,
keyPrefix)), pattern);
URI expected = URI.create(StringUtils.format("gs://%s/%s", bucket,
storageObject3.getName()));
@@ -70,7 +70,7 @@ public class GoogleTimestampVersionedDataFinderTest
storageObject4.setLastUpdateTimeMillis(System.currentTimeMillis() + 100);
final GoogleStorage storage =
ObjectStorageIteratorTest.makeMockClient(ImmutableList.of(storageObject1,
storageObject2, storageObject3, storageObject4));
- final GoogleTimestampVersionedDataFinder finder = new
GoogleTimestampVersionedDataFinder(storage);
+ final GoogleTimestampVersionedDataFinder finder = new
GoogleTimestampVersionedDataFinder(storage, new GoogleInputDataConfig());
Pattern pattern = Pattern.compile("v.*");
URI latest =
finder.getLatestVersion(URI.create(StringUtils.format("gs://%s/%s", bucket,
keyPrefix)), pattern);
URI expected = URI.create(StringUtils.format("gs://%s/%s", bucket,
storageObject3.getName()));
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]