Abacn commented on code in PR #40279:
URL: https://github.com/apache/beam/pull/40279#discussion_r4136467406
##########
sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java:
##########
Review Comment:
"thrown" is now unused.
##########
sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilIT.java:
##########
@@ -53,12 +93,778 @@
* <p>This is a runnerless integration test, even though the Beam IT framework
assumes one. Thus,
* this test should only be run against single runner (such as DirectRunner).
*/
-@RunWith(JUnit4.class)
+@RunWith(Parameterized.class)
@Category(UsesKms.class)
public class GcsUtilIT {
+
+ private static final String READ_COUNTER_PREFIX = "it_read_bytes";
+ private static final String WRITE_COUNTER_PREFIX = "it_write_bytes";
+
+ @Parameters(name = "{0}")
+ public static Iterable<String> data() {
+ return Arrays.asList("use_gcsutil_v1", "use_gcsutil_v2");
+ }
+
+ @Parameter public String experiment;
+
+ private TestPipelineOptions options;
+ private GcsUtil gcsUtil;
+
+ @Before
+ public void setUp() {
+ options =
TestPipeline.testingPipelineOptions().as(TestPipelineOptions.class);
+
+ // set the experimental flag.
+ ExperimentalOptions experimentalOptions =
options.as(ExperimentalOptions.class);
+ experimentalOptions.setExperiments(Collections.singletonList(experiment));
+
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ gcsUtil = gcsOptions.getGcsUtil();
+ }
+
+ /** Returns a bucket name unique to this test run, so concurrent runs don't
collide. */
+ private static String randomBucketName() {
+ return "apache-beam-temp-bucket-" + UUID.randomUUID();
Review Comment:
We now randomize bucket names. It resolves conflict, however, there is a
risk of bucket leak. As a follow up consider add a gcs temp bucket cleaner in
https://github.com/apache/beam/tree/master/.test-infra/tools and as part of
cleanUp workflow
##########
sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java:
##########
@@ -427,4 +443,128 @@ public void testDefaultUploadChunkSizeMatchesV1() {
AsyncWriteChannelOptions.DEFAULT.getUploadChunkSize(),
GcsUtilV2.DEFAULT_UPLOAD_CHUNK_SIZE_BYTES);
}
+
+ //
---------------------------------------------------------------------------------------------
+ // Behavior shared with GcsUtilV1, through a mocked java-storage client.
Each test mirrors the
+ // GcsUtilV1Test case it names; end-to-end parity is covered by GcsUtilIT.
+ //
---------------------------------------------------------------------------------------------
+
+ /**
+ * Returns a {@link GcsUtil} backed by a real {@link GcsUtilV2} that issues
every call to {@code
+ * storage}. Performance metrics are off by default, so the per-operation
clients of {@link
+ * GcsUtilV2#storageWithHttpMetrics} resolve to this one as well.
+ */
+ private GcsUtil gcsUtilWithV2Storage(com.google.cloud.storage.Storage
storage) {
Review Comment:
nit: trim fully qualified names
##########
sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilIT.java:
##########
@@ -53,12 +93,778 @@
* <p>This is a runnerless integration test, even though the Beam IT framework
assumes one. Thus,
* this test should only be run against single runner (such as DirectRunner).
*/
-@RunWith(JUnit4.class)
+@RunWith(Parameterized.class)
@Category(UsesKms.class)
public class GcsUtilIT {
+
+ private static final String READ_COUNTER_PREFIX = "it_read_bytes";
+ private static final String WRITE_COUNTER_PREFIX = "it_write_bytes";
+
+ @Parameters(name = "{0}")
+ public static Iterable<String> data() {
+ return Arrays.asList("use_gcsutil_v1", "use_gcsutil_v2");
+ }
+
+ @Parameter public String experiment;
+
+ private TestPipelineOptions options;
+ private GcsUtil gcsUtil;
+
+ @Before
+ public void setUp() {
+ options =
TestPipeline.testingPipelineOptions().as(TestPipelineOptions.class);
+
+ // set the experimental flag.
+ ExperimentalOptions experimentalOptions =
options.as(ExperimentalOptions.class);
+ experimentalOptions.setExperiments(Collections.singletonList(experiment));
+
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ gcsUtil = gcsOptions.getGcsUtil();
+ }
+
+ /** Returns a bucket name unique to this test run, so concurrent runs don't
collide. */
+ private static String randomBucketName() {
+ return "apache-beam-temp-bucket-" + UUID.randomUUID();
+ }
+
+ @Test
+ public void testFileSize() throws IOException {
+ final GcsPath gcsPath =
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kinglear.txt");
+ final long expectedSize = 157283L;
+
+ assertEquals(expectedSize, gcsUtil.fileSize(gcsPath));
+ }
+
+ @Test
+ public void testGetObjectOrGetBlob() throws IOException {
+ final GcsPath existingPath =
+ GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kinglear.txt");
+ final String expectedCRC = "s0a3Tg==";
+
+ String crc;
+ if (experiment.equals("use_gcsutil_v2")) {
+ Blob blob = gcsUtil.getBlob(existingPath);
+ crc = blob.getCrc32c();
+ } else {
+ StorageObject obj = gcsUtil.getObject(existingPath);
+ crc = obj.getCrc32c();
+ }
+ assertEquals(expectedCRC, crc);
+
+ final GcsPath nonExistentPath =
+ GcsPath.fromUri("gs://my-random-test-bucket-12345/unknown-12345.txt");
+ final GcsPath forbiddenPath =
GcsPath.fromUri("gs://test-bucket/unknown-12345.txt");
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.getBlob(nonExistentPath));
+ // For V2, we are returning AccessDeniedException (a subclass of
IOException) for forbidden
+ // paths.
+ assertThrows(AccessDeniedException.class, () ->
gcsUtil.getBlob(forbiddenPath));
+ } else {
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.getObject(nonExistentPath));
+ assertThrows(IOException.class, () -> gcsUtil.getObject(forbiddenPath));
+ }
+ }
+
+ @Test
+ public void testGetObjectsOrGetBlobs() throws IOException {
+ final GcsPath existingPath =
+ GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kinglear.txt");
+ final GcsPath nonExistentPath =
+ GcsPath.fromUri("gs://my-random-test-bucket-12345/unknown-12345.txt");
+ final List<GcsPath> paths = Arrays.asList(existingPath, nonExistentPath);
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ List<GcsUtilV2.BlobResult> results = gcsUtil.getBlobs(paths);
+ assertEquals(2, results.size());
+ assertTrue(results.get(0).blob() != null);
+ assertTrue(results.get(0).ioException() == null);
+ assertTrue(results.get(1).blob() == null);
+ assertTrue(results.get(1).ioException() != null);
+ } else {
+ List<GcsUtil.StorageObjectOrIOException> results =
gcsUtil.getObjects(paths);
+ assertEquals(2, results.size());
+ assertTrue(results.get(0).storageObject() != null);
+ assertTrue(results.get(0).ioException() == null);
+ assertTrue(results.get(1).storageObject() == null);
+ assertTrue(results.get(1).ioException() != null);
+ }
+ }
+
+ @Test
+ public void testListObjectsOrListBlobs() throws IOException {
+ final String bucket = "apache-beam-samples";
+ final String prefix = "shakespeare/kingrichard";
+
+ List<String> names;
+ if (experiment.equals("use_gcsutil_v2")) {
+ Page<Blob> blobs = gcsUtil.listBlobs(bucket, prefix, null);
+ names = blobs.streamAll().map(blob ->
blob.getName()).collect(Collectors.toList());
+ } else {
+ Objects objs = gcsUtil.listObjects(bucket, prefix, null);
+ names = objs.getItems().stream().map(obj ->
obj.getName()).collect(Collectors.toList());
+ }
+ assertEquals(
+ Arrays.asList("shakespeare/kingrichardii.txt",
"shakespeare/kingrichardiii.txt"), names);
+
+ final String randomPrefix = "my-random-prefix/random";
+ if (experiment.equals("use_gcsutil_v2")) {
+ Page<Blob> blobs = gcsUtil.listBlobs(bucket, randomPrefix, null);
+ assertEquals(0, blobs.streamAll().count());
+ } else {
+ Objects objs = gcsUtil.listObjects(bucket, randomPrefix, null);
+ assertEquals(null, objs.getItems());
+ }
+ }
+
+ @Test
+ public void testExpand() throws IOException {
+ final GcsPath existingPattern =
+
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kingrichardii*.txt");
+ List<GcsPath> paths = gcsUtil.expand(existingPattern);
+
+ assertEquals(
+ Arrays.asList(
+
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kingrichardii.txt"),
+
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kingrichardiii.txt")),
+ paths);
+
+ final GcsPath nonExistentPattern1 =
+
GcsPath.fromUri("gs://apache-beam-samples/my_random_folder/random*.txt");
+ assertTrue(gcsUtil.expand(nonExistentPattern1).isEmpty());
+
+ final GcsPath nonExistentPattern2 =
+ GcsPath.fromUri("gs://apache-beam-samples/shakespeare/king*.csv");
+ assertTrue(gcsUtil.expand(nonExistentPattern2).isEmpty());
+ }
+
+ @Test
+ public void testGetBucketOrGetBucketWithOptions() throws IOException {
+ final GcsPath existingPath = GcsPath.fromUri("gs://apache-beam-samples");
+
+ String bucket;
+ if (experiment.equals("use_gcsutil_v2")) {
+ bucket = gcsUtil.getBucketWithOptions(existingPath).getName();
+ } else {
+ bucket = gcsUtil.getBucket(existingPath).getName();
+ }
+ assertEquals("apache-beam-samples", bucket);
+
+ final GcsPath nonExistentPath =
GcsPath.fromUri("gs://my-random-test-bucket-12345");
+ final GcsPath forbiddenPath = GcsPath.fromUri("gs://test-bucket");
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ assertThrows(
+ FileNotFoundException.class, () ->
gcsUtil.getBucketWithOptions(nonExistentPath));
+ assertThrows(AccessDeniedException.class, () ->
gcsUtil.getBucketWithOptions(forbiddenPath));
+ } else {
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.getBucket(nonExistentPath));
+ assertThrows(AccessDeniedException.class, () ->
gcsUtil.getBucket(forbiddenPath));
+ }
+ }
+
+ @Test
+ public void testBucketAccessible() throws IOException {
+ final GcsPath existingPath = GcsPath.fromUri("gs://apache-beam-samples");
+ final GcsPath nonExistentPath =
GcsPath.fromUri("gs://my-random-test-bucket-12345");
+ final GcsPath forbiddenPath = GcsPath.fromUri("gs://test-bucket");
+
+ assertEquals(true, gcsUtil.bucketAccessible(existingPath));
+ assertEquals(false, gcsUtil.bucketAccessible(nonExistentPath));
+ assertEquals(false, gcsUtil.bucketAccessible(forbiddenPath));
+ }
+
+ @Test
+ public void testBucketOwner() throws IOException {
+ final GcsPath existingPath = GcsPath.fromUri("gs://apache-beam-samples");
+ final long expectedProjectNumber = 844138762903L; // apache-beam-testing
+ assertEquals(expectedProjectNumber, gcsUtil.bucketOwner(existingPath));
+
+ final GcsPath nonExistentPath =
GcsPath.fromUri("gs://my-random-test-bucket-12345");
+ final GcsPath forbiddenPath = GcsPath.fromUri("gs://test-bucket");
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.bucketOwner(nonExistentPath));
+ assertThrows(AccessDeniedException.class, () ->
gcsUtil.bucketOwner(forbiddenPath));
+ }
+
+ @Test
+ public void testCreateAndRemoveBucket() throws IOException {
+ final GcsPath gcsPath = GcsPath.fromUri("gs://" + randomBucketName());
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ BucketInfo bucketInfo = BucketInfo.of(gcsPath.getBucket());
+ try {
+ assertFalse(gcsUtil.bucketAccessible(gcsPath));
+ gcsUtil.createBucket(bucketInfo);
+ assertTrue(gcsUtil.bucketAccessible(gcsPath));
+
+ // raise exception when the bucket already exists during creation
+ assertThrows(FileAlreadyExistsException.class, () ->
gcsUtil.createBucket(bucketInfo));
+
+ assertTrue(gcsUtil.bucketAccessible(gcsPath));
+ gcsUtil.removeBucket(bucketInfo);
+ assertFalse(gcsUtil.bucketAccessible(gcsPath));
+
+ // raise exception when the bucket does not exist during removal
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.removeBucket(bucketInfo));
+ } finally {
+ // clean up and ignore errors no matter what
+ try {
+ gcsUtil.removeBucket(bucketInfo);
+ } catch (IOException e) {
+ }
+ }
+ } else {
+ Bucket bucket = new Bucket().setName(gcsPath.getBucket());
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ String projectId = gcsOptions.getProject();
+ try {
+ assertFalse(gcsUtil.bucketAccessible(gcsPath));
+ gcsUtil.createBucket(projectId, bucket);
+ assertTrue(gcsUtil.bucketAccessible(gcsPath));
+
+ // raise exception when the bucket already exists during creation
+ assertThrows(
+ FileAlreadyExistsException.class, () ->
gcsUtil.createBucket(projectId, bucket));
+
+ assertTrue(gcsUtil.bucketAccessible(gcsPath));
+ gcsUtil.removeBucket(bucket);
+ assertFalse(gcsUtil.bucketAccessible(gcsPath));
+
+ // raise exception when the bucket does not exist during removal
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.removeBucket(bucket));
+ } finally {
+ // clean up and ignore errors no matter what
+ try {
+ gcsUtil.removeBucket(bucket);
+ } catch (IOException e) {
+ }
+ }
+ }
+ }
+
+ private List<GcsPath> createTestBucketHelper(String bucketName, boolean
copyData)
+ throws IOException {
+ final List<GcsPath> originPaths =
+ Arrays.asList(
+
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kingrichardii.txt"),
+
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kingrichardiii.txt"));
+
+ final List<GcsPath> testPaths =
+ originPaths.stream()
+ .map(o -> GcsPath.fromComponents(bucketName, o.getObject()))
+ .collect(Collectors.toList());
+
+ // create bucket and copy some initial files into there
+ if (experiment.equals("use_gcsutil_v2")) {
+ gcsUtil.createBucket(BucketInfo.of(bucketName));
+
+ if (copyData) {
+ gcsUtil.copyV2(originPaths, testPaths);
+ } else {
+ return Collections.emptyList();
+ }
+ } else {
+ GcsOptions gcsOptions = options.as(GcsOptions.class);
+ gcsUtil.createBucket(gcsOptions.getProject(), new
Bucket().setName(bucketName));
+
+ if (copyData) {
+ final List<String> originList =
+ originPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> testList =
+ testPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ gcsUtil.copy(originList, testList);
+ } else {
+ return Collections.emptyList();
+ }
+ }
+
+ return testPaths;
+ }
+
+ private void tearDownTestBucketHelper(String bucketName) {
+ try {
+ // use "**" in the pattern to match any characters including "/".
+ final List<GcsPath> paths =
+ gcsUtil.expand(GcsPath.fromUri(String.format("gs://%s/**",
bucketName)));
+ if (experiment.equals("use_gcsutil_v2")) {
+ gcsUtil.remove(paths, MissingStrategy.SKIP_IF_MISSING);
+ gcsUtil.removeBucket(BucketInfo.of(bucketName));
+ } else {
+
gcsUtil.remove(paths.stream().map(GcsPath::toString).collect(Collectors.toList()));
+ gcsUtil.removeBucket(new Bucket().setName(bucketName));
+ }
+ } catch (IOException e) {
+ System.err.println(
+ "Error during tear down of test bucket " + bucketName + ": " +
e.getMessage());
+ }
+ }
+
+ @Test
+ public void testCopy() throws IOException {
+ final String existingBucket = randomBucketName();
+ final String nonExistentBucket = "my-random-test-bucket-12345";
+
+ try {
+ final List<GcsPath> srcPaths = createTestBucketHelper(existingBucket,
true);
+ final List<GcsPath> dstPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(existingBucket, o.getObject() +
".bak"))
+ .collect(Collectors.toList());
+ final List<GcsPath> errPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(nonExistentBucket,
o.getObject()))
+ .collect(Collectors.toList());
+
+ assertNotExists(dstPaths.get(0));
+ assertNotExists(dstPaths.get(1));
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ // (1) when the target files do not exist
+ gcsUtil.copyV2(srcPaths, dstPaths);
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+
+ // (2) when the target files exist
+ // (2a) no exception on SAFE_OVERWRITE, ALWAYS_OVERWRITE,
SKIP_IF_EXISTS
+ gcsUtil.copyV2(srcPaths, dstPaths);
+ gcsUtil.copy(srcPaths, dstPaths, OverwriteStrategy.ALWAYS_OVERWRITE);
+ gcsUtil.copy(srcPaths, dstPaths, OverwriteStrategy.SKIP_IF_EXISTS);
+
+ // (2b) raise exception on FAIL_IF_EXISTS
+ assertThrows(
+ FileAlreadyExistsException.class,
+ () -> gcsUtil.copy(srcPaths, dstPaths,
OverwriteStrategy.FAIL_IF_EXISTS));
+
+ // (3) raise exception when the target bucket is nonexistent.
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.copyV2(srcPaths, errPaths));
+
+ // (4) raise exception when the source files are nonexistent.
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.copyV2(errPaths, dstPaths));
+ } else {
+ final List<String> srcList =
+ srcPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> dstList =
+ dstPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> errList =
+ errPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+
+ // (1) when the target files do not exist
+ gcsUtil.copy(srcList, dstList);
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+
+ // (2) when the target files exist, no exception
+ gcsUtil.copy(srcList, dstList);
+
+ // (3) raise exception when the target bucket is nonexistent.
+ assertThrows(FileNotFoundException.class, () -> gcsUtil.copy(srcList,
errList));
+
+ // (4) raise exception when the source files are nonexistent.
+ assertThrows(FileNotFoundException.class, () -> gcsUtil.copy(errList,
dstList));
+ }
+ } finally {
+ tearDownTestBucketHelper(existingBucket);
+ }
+ }
+
+ @Test
+ public void testRemove() throws IOException {
+ final String existingBucket = randomBucketName();
+ final String nonExistentBucket = "my-random-test-bucket-12345";
+
+ try {
+ final List<GcsPath> srcPaths = createTestBucketHelper(existingBucket,
true);
+ final List<GcsPath> errPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(nonExistentBucket,
o.getObject()))
+ .collect(Collectors.toList());
+
+ assertExists(srcPaths.get(0));
+ assertExists(srcPaths.get(1));
+
+ if (experiment.equals("use_gcsutil_v2")) {
+ // (1) when the files to remove exist
+ gcsUtil.removeV2(srcPaths);
+ assertNotExists(srcPaths.get(0));
+ assertNotExists(srcPaths.get(1));
+
+ // (2) when the files to remove have been deleted
+ // (2a) no exception on SKIP_IF_MISSING
+ gcsUtil.removeV2(srcPaths);
+ gcsUtil.remove(srcPaths, MissingStrategy.SKIP_IF_MISSING);
+
+ // (2b) raise exception on FAIL_IF_MISSING
+ assertThrows(
+ FileNotFoundException.class,
+ () -> gcsUtil.remove(srcPaths, MissingStrategy.FAIL_IF_MISSING));
+
+ // (3) when the files are from an nonexistent bucket
+ // (3a) no exception on SKIP_IF_MISSING
+ gcsUtil.removeV2(errPaths);
+ gcsUtil.remove(errPaths, MissingStrategy.SKIP_IF_MISSING);
+
+ // (3b) raise exception on FAIL_IF_MISSING
+ assertThrows(
+ FileNotFoundException.class,
+ () -> gcsUtil.remove(errPaths, MissingStrategy.FAIL_IF_MISSING));
+ } else {
+ final List<String> srcList =
+ srcPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> errList =
+ errPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+
+ // (1) when the files to remove exist
+ gcsUtil.remove(srcList);
+ assertNotExists(srcPaths.get(0));
+ assertNotExists(srcPaths.get(1));
+
+ // (2) when the files to remove have been deleted, no exception
+ gcsUtil.remove(srcList);
+
+ // (3) when the files are from an nonexistent bucket, no exception
+ gcsUtil.remove(errList);
+ }
+ } finally {
+ tearDownTestBucketHelper(existingBucket);
+ }
+ }
+
+ @Test
+ public void testRename() throws IOException {
+ final String existingBucket = randomBucketName();
+ final String nonExistentBucket = "my-random-test-bucket-12345";
+
+ try {
+ final List<GcsPath> srcPaths = createTestBucketHelper(existingBucket,
true);
+ final List<GcsPath> tmpPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(existingBucket, "tmp/" +
o.getObject()))
+ .collect(Collectors.toList());
+ final List<GcsPath> dstPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(existingBucket, o.getObject() +
".bak"))
+ .collect(Collectors.toList());
+ final List<GcsPath> errPaths =
+ srcPaths.stream()
+ .map(o -> GcsPath.fromComponents(nonExistentBucket,
o.getObject()))
+ .collect(Collectors.toList());
+
+ assertNotExists(dstPaths.get(0));
+ assertNotExists(dstPaths.get(1));
+ if (experiment.equals("use_gcsutil_v2")) {
+ // Make a copy of sources
+ gcsUtil.copyV2(srcPaths, tmpPaths);
+
+ // (1) when the source files exist and target files do not
+ gcsUtil.renameV2(tmpPaths, dstPaths);
+ assertNotExists(tmpPaths.get(0));
+ assertNotExists(tmpPaths.get(1));
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+
+ // (2) when the source files do not exist
+ // (2a) no exception if IGNORE_MISSING_FILES is set
+ gcsUtil.renameV2(errPaths, dstPaths,
MoveOptions.StandardMoveOptions.IGNORE_MISSING_FILES);
+
+ // (2b) raise exception if if IGNORE_MISSING_FILES is not set
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.renameV2(errPaths, dstPaths));
+
+ // (3) when both source files and target files exist
+ gcsUtil.renameV2(
+ srcPaths, dstPaths,
MoveOptions.StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS);
+ gcsUtil.renameV2(srcPaths, dstPaths);
+ } else {
+ final List<String> srcList =
+ srcPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> tmpList =
+ tmpPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> dstList =
+ dstPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+ final List<String> errList =
+ errPaths.stream().map(o ->
o.toString()).collect(Collectors.toList());
+
+ // Make a copy of sources
+ gcsUtil.copy(srcList, tmpList);
+
+ // (1) when the source files exist and target files do not
+ gcsUtil.rename(tmpList, dstList);
+ assertNotExists(tmpPaths.get(0));
+ assertNotExists(tmpPaths.get(1));
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+
+ // (2) when the source files do not exist
+ // (2a) no exception if IGNORE_MISSING_FILES is set
+ gcsUtil.rename(errList, dstList,
MoveOptions.StandardMoveOptions.IGNORE_MISSING_FILES);
+
+ // (2b) raise exception if if IGNORE_MISSING_FILES is not set
+ assertThrows(FileNotFoundException.class, () ->
gcsUtil.rename(errList, dstList));
+
+ // (3) when both source files and target files exist
+ assertExists(srcPaths.get(0));
+ assertExists(srcPaths.get(1));
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+
+ // There is a bug in V1 where SKIP_IF_DESTINATION_EXISTS is not
honored.
+ gcsUtil.rename(
+ srcList, dstList,
MoveOptions.StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS);
+
+ assertNotExists(srcPaths.get(0)); // BUG! The renaming is supposed to
be skipped
+ assertNotExists(srcPaths.get(1)); // BUG! The renaming is supposed to
be skipped
+ // assertExists(srcPaths.get(0));
+ // assertExists(srcPaths.get(1));
+ assertExists(dstPaths.get(0));
+ assertExists(dstPaths.get(1));
+ }
+ } finally {
+ tearDownTestBucketHelper(existingBucket);
+ }
+ }
+
+ private void assertExists(GcsPath path) throws IOException {
+ if (experiment.equals("use_gcsutil_v2")) {
+ gcsUtil.getBlob(path);
+ } else {
+ gcsUtil.getObject(path);
+ }
+ }
+
+ private void assertNotExists(GcsPath path) throws IOException {
+ if (experiment.equals("use_gcsutil_v2")) {
+ assertThrows(FileNotFoundException.class, () -> gcsUtil.getBlob(path));
+ } else {
+ assertThrows(FileNotFoundException.class, () -> gcsUtil.getObject(path));
+ }
+ }
+
+ String computeHash(ByteBuffer buffer) throws NoSuchAlgorithmException {
+ MessageDigest digest = MessageDigest.getInstance("SHA-256");
+ digest.update(buffer);
+ byte[] hashBytes = digest.digest();
+
+ // Convert bytes to Hex String
+ StringBuilder sb = new StringBuilder();
+ for (byte b : hashBytes) {
+ sb.append(String.format("%02x", b));
+ }
+ return sb.toString();
+ }
+
+ @Test
+ public void testRead() throws IOException, NoSuchAlgorithmException {
+ final GcsPath gcsPath =
GcsPath.fromUri("gs://apache-beam-samples/shakespeare/kinglear.txt");
+ final String expectedHash =
"674a2725884307c96398440497c889ad8cecccedf5689df85e6b0faabe4e0fe8";
+ final long expectedSize = 157283L;
+
+ try (SeekableByteChannel channel = gcsUtil.open(gcsPath)) {
+ // Verify Size
+ assertEquals(expectedSize, channel.size());
+ assertEquals(0, channel.position());
+
+ // Read content into ByteBuffer.
+ // Allocate a larger buffer to ensure we receive the EOF at the expected
place.
+ ByteBuffer buffer = ByteBuffer.allocate((int) expectedSize + 1024);
+ int bytesRead = StorageChannelUtils.blockingFillFrom(buffer, channel);
+
+ // Verify total bytes read and position
+ assertEquals(expectedSize, bytesRead);
+ assertEquals(expectedSize, channel.position());
+
+ // Flip the buffer to prepare it for reading (sets limit to current
position, position to 0)
+ buffer.flip();
+
+ // Verify hash
+ String actualHash = computeHash(buffer);
+ assertEquals("Content hash should match", expectedHash, actualHash);
+ }
+ }
+
+ @Test
+ public void testWriteAndRead() throws IOException {
+ final String bucketName = randomBucketName();
+ final GcsPath targetPath =
+ GcsPath.fromComponents(bucketName, "test-object-" +
java.util.UUID.randomUUID() + ".txt");
Review Comment:
nit: trim fully qualified name
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]