This is an automated email from the ASF dual-hosted git repository.
pvillard pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/master by this push:
new a700890 NIFI-7298: PutAzureDataLakeStorage tests.
a700890 is described below
commit a7008903edf3cc751a48695ba09701667892e545
Author: Peter Turcsanyi <[email protected]>
AuthorDate: Wed Apr 22 23:22:04 2020 +0200
NIFI-7298: PutAzureDataLakeStorage tests.
Signed-off-by: Pierre Villard <[email protected]>
This closes #4227.
---
.../azure/storage/PutAzureDataLakeStorage.java | 41 ++--
.../storage/AbstractAzureDataLakeStorageIT.java | 57 +++++
.../azure/storage/AbstractAzureStorageIT.java | 7 +-
.../storage/ITDeleteAzureDataLakeStorage.java | 2 +
.../azure/storage/ITFetchAzureDataLakeStorage.java | 2 +
.../azure/storage/ITPutAzureDataLakeStorage.java | 235 +++++++++++++++++++--
.../storage/TestAbstractAzureDataLakeStorage.java | 124 +++++++++++
7 files changed, 437 insertions(+), 31 deletions(-)
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
index b59c366..94fa90f 100644
---
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/main/java/org/apache/nifi/processors/azure/storage/PutAzureDataLakeStorage.java
@@ -16,12 +16,11 @@
*/
package org.apache.nifi.processors.azure.storage;
-import java.io.BufferedInputStream;
-import java.io.InputStream;
-import java.util.HashMap;
-import java.util.Map;
-import java.util.concurrent.TimeUnit;
-
+import com.azure.storage.file.datalake.DataLakeDirectoryClient;
+import com.azure.storage.file.datalake.DataLakeFileClient;
+import com.azure.storage.file.datalake.DataLakeFileSystemClient;
+import com.azure.storage.file.datalake.DataLakeServiceClient;
+import org.apache.commons.lang3.StringUtils;
import org.apache.nifi.annotation.behavior.InputRequirement;
import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
import org.apache.nifi.annotation.behavior.WritesAttribute;
@@ -33,13 +32,13 @@ import org.apache.nifi.flowfile.FlowFile;
import org.apache.nifi.processor.ProcessContext;
import org.apache.nifi.processor.ProcessSession;
import org.apache.nifi.processor.exception.ProcessException;
-
-import com.azure.storage.file.datalake.DataLakeDirectoryClient;
-import com.azure.storage.file.datalake.DataLakeFileClient;
-import com.azure.storage.file.datalake.DataLakeFileSystemClient;
-import com.azure.storage.file.datalake.DataLakeServiceClient;
import org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
+import java.io.BufferedInputStream;
+import java.io.InputStream;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
@Tags({"azure", "microsoft", "cloud", "storage", "adlsgen2", "datalake"})
@SeeAlso({DeleteAzureDataLakeStorage.class, FetchAzureDataLakeStorage.class})
@@ -50,7 +49,6 @@ import
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
@WritesAttribute(attribute = "azure.primaryUri", description =
"Primary location for file content"),
@WritesAttribute(attribute = "azure.length", description = "Length of
the file")})
@InputRequirement(Requirement.INPUT_REQUIRED)
-
public class PutAzureDataLakeStorage extends
AbstractAzureDataLakeStorageProcessor {
@Override
@@ -59,23 +57,35 @@ public class PutAzureDataLakeStorage extends
AbstractAzureDataLakeStorageProcess
if (flowFile == null) {
return;
}
+
final long startNanos = System.nanoTime();
try {
final String fileSystem =
context.getProperty(FILESYSTEM).evaluateAttributeExpressions(flowFile).getValue();
final String directory =
context.getProperty(DIRECTORY).evaluateAttributeExpressions(flowFile).getValue();
final String fileName =
context.getProperty(FILE).evaluateAttributeExpressions(flowFile).getValue();
+
+ if (StringUtils.isBlank(fileSystem)) {
+ throw new ProcessException(FILESYSTEM.getDisplayName() + "
property evaluated to empty string. " +
+ FILESYSTEM.getDisplayName() + " must be specified as a
non-empty string.");
+ }
+ if (StringUtils.isBlank(fileName)) {
+ throw new ProcessException(FILE.getDisplayName() + " property
evaluated to empty string. " +
+ FILE.getDisplayName() + " must be specified as a
non-empty string.");
+ }
+
final DataLakeServiceClient storageClient =
getStorageClient(context, flowFile);
- final DataLakeFileSystemClient dataLakeFileSystemClient =
storageClient.getFileSystemClient(fileSystem);
- final DataLakeDirectoryClient directoryClient =
dataLakeFileSystemClient.getDirectoryClient(directory);
+ final DataLakeFileSystemClient fileSystemClient =
storageClient.getFileSystemClient(fileSystem);
+ final DataLakeDirectoryClient directoryClient =
fileSystemClient.getDirectoryClient(directory);
final DataLakeFileClient fileClient =
directoryClient.createFile(fileName);
+
final long length = flowFile.getSize();
if (length > 0) {
try (final InputStream rawIn = session.read(flowFile); final
BufferedInputStream in = new BufferedInputStream(rawIn)) {
fileClient.append(in, 0, length);
-
}
}
fileClient.flush(length);
+
final Map<String, String> attributes = new HashMap<>();
attributes.put("azure.filesystem", fileSystem);
attributes.put("azure.directory", directory);
@@ -84,7 +94,6 @@ public class PutAzureDataLakeStorage extends
AbstractAzureDataLakeStorageProcess
attributes.put("azure.length", String.valueOf(length));
flowFile = session.putAllAttributes(flowFile, attributes);
-
session.transfer(flowFile, REL_SUCCESS);
final long transferMillis =
TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos);
session.getProvenanceReporter().send(flowFile,
fileClient.getFileUrl(), transferMillis);
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
new file mode 100644
index 0000000..a9b2bd9
--- /dev/null
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureDataLakeStorageIT.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.processors.azure.storage;
+
+import com.azure.storage.common.StorageSharedKeyCredential;
+import com.azure.storage.file.datalake.DataLakeFileSystemClient;
+import com.azure.storage.file.datalake.DataLakeServiceClient;
+import com.azure.storage.file.datalake.DataLakeServiceClientBuilder;
+import org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor;
+import org.junit.After;
+import org.junit.Before;
+
+import java.util.UUID;
+
+public abstract class AbstractAzureDataLakeStorageIT extends
AbstractAzureStorageIT {
+
+ private static final String FILESYSTEM_NAME_PREFIX =
"nifi-test-filesystem";
+
+ protected String fileSystemName;
+ protected DataLakeFileSystemClient fileSystemClient;
+
+ @Before
+ public void setUpAzureDataLakeStorageIT() {
+ fileSystemName = String.format("%s-%s", FILESYSTEM_NAME_PREFIX,
UUID.randomUUID());
+
+ runner.setProperty(AbstractAzureDataLakeStorageProcessor.FILESYSTEM,
fileSystemName);
+
+ DataLakeServiceClient storageClient = createStorageClient();
+ fileSystemClient = storageClient.createFileSystem(fileSystemName);
+ }
+
+ @After
+ public void tearDownAzureDataLakeStorageIT() {
+ fileSystemClient.delete();
+ }
+
+ private DataLakeServiceClient createStorageClient() {
+ return new DataLakeServiceClientBuilder()
+ .endpoint("https://" + getAccountName() +
".dfs.core.windows.net")
+ .credential(new StorageSharedKeyCredential(getAccountName(),
getAccountKey()))
+ .buildClient();
+ }
+}
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
index f5fa2e5..7689b96 100644
---
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/AbstractAzureStorageIT.java
@@ -35,15 +35,16 @@ import java.util.Properties;
import static org.junit.Assert.fail;
-public abstract class AbstractAzureStorageIT {private static final Properties
CONFIG;
+public abstract class AbstractAzureStorageIT {
+
+ private static final Properties CONFIG;
private static final String CREDENTIALS_FILE =
System.getProperty("user.home") + "/azure-credentials.PROPERTIES";
static {
- final FileInputStream fis;
CONFIG = new Properties();
try {
- fis = new FileInputStream(CREDENTIALS_FILE);
+ final FileInputStream fis = new FileInputStream(CREDENTIALS_FILE);
try {
CONFIG.load(fis);
} catch (IOException e) {
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
index e8c0230..29915eb 100644
---
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITDeleteAzureDataLakeStorage.java
@@ -19,6 +19,7 @@ package org.apache.nifi.processors.azure.storage;
import org.apache.nifi.processor.Processor;
import org.apache.nifi.util.MockFlowFile;
import org.junit.Before;
+import org.junit.Ignore;
import org.junit.Test;
import java.util.List;
@@ -35,6 +36,7 @@ public class ITDeleteAzureDataLakeStorage extends
AbstractAzureBlobStorageIT {
runner.setProperty(DeleteAzureDataLakeStorage.FILE, TEST_FILE_NAME);
}
+ @Ignore
@Test
public void testDeleteFile() throws Exception {
runner.assertValid();
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
index 93a414f..a242ede 100644
---
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITFetchAzureDataLakeStorage.java
@@ -19,6 +19,7 @@ package org.apache.nifi.processors.azure.storage;
import org.apache.nifi.processor.Processor;
import org.apache.nifi.util.MockFlowFile;
import org.junit.Before;
+import org.junit.Ignore;
import org.junit.Test;
import java.util.List;
@@ -36,6 +37,7 @@ public class ITFetchAzureDataLakeStorage extends
AbstractAzureBlobStorageIT {
}
+ @Ignore
@Test
public void testFetchFile() throws Exception {
runner.assertValid();
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
index dba8abd..1a6105f 100644
---
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/ITPutAzureDataLakeStorage.java
@@ -16,14 +16,32 @@
*/
package org.apache.nifi.processors.azure.storage;
+import com.azure.storage.file.datalake.DataLakeDirectoryClient;
+import com.azure.storage.file.datalake.DataLakeFileClient;
+import com.google.common.net.UrlEscapers;
+import org.apache.commons.lang3.StringUtils;
import org.apache.nifi.processor.Processor;
import org.apache.nifi.util.MockFlowFile;
import org.junit.Before;
+import org.junit.Ignore;
import org.junit.Test;
-import java.util.List;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
-public class ITPutAzureDataLakeStorage extends AbstractAzureBlobStorageIT {
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+public class ITPutAzureDataLakeStorage extends AbstractAzureDataLakeStorageIT {
+
+ private static final String DIRECTORY = "dir1";
+ private static final String FILE_NAME = "file1";
+ private static final byte[] FILE_DATA = "0123456789".getBytes();
+
+ private static final String EL_FILESYSTEM = "az.filesystem";
+ private static final String EL_DIRECTORY = "az.directory";
+ private static final String EL_FILE_NAME = "az.filename";
@Override
protected Class<? extends Processor> getProcessorClass() {
@@ -32,26 +50,219 @@ public class ITPutAzureDataLakeStorage extends
AbstractAzureBlobStorageIT {
@Before
public void setUp() {
- runner.setProperty(PutAzureDataLakeStorage.FILE, TEST_FILE_NAME);
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, DIRECTORY);
+ runner.setProperty(PutAzureDataLakeStorage.FILE, FILE_NAME);
+ }
+
+ @Test
+ public void testPutFileToExistingDirectory() throws Exception {
+ fileSystemClient.createDirectory(DIRECTORY);
+
+ runProcessor(FILE_DATA);
+
+ assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+ }
+
+ @Test
+ public void testPutFileToNonExistingDirectory() throws Exception {
+ runProcessor(FILE_DATA);
+
+ assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+ }
+
+ @Test
+ public void testPutFileToDeepDirectory() throws Exception {
+ String baseDirectory = "dir1/dir2";
+ String fullDirectory = baseDirectory + "/dir3/dir4";
+ fileSystemClient.createDirectory(baseDirectory);
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, fullDirectory);
+
+ runProcessor(FILE_DATA);
+
+ assertSuccess(fullDirectory, FILE_NAME, FILE_DATA);
+ }
+
+ @Test
+ public void testPutFileToRootDirectory() throws Exception {
+ String rootDirectory = "";
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, rootDirectory);
+
+ runProcessor(FILE_DATA);
+
+ assertSuccess(rootDirectory, FILE_NAME, FILE_DATA);
+ }
+
+ @Test
+ public void testPutEmptyFile() throws Exception {
+ byte[] fileData = new byte[0];
+
+ runProcessor(fileData);
+
+ assertSuccess(DIRECTORY, FILE_NAME, fileData);
+ }
+
+ @Ignore
+ // ignore excessive test with larger file size
+ @Test
+ public void testPutBigFile() throws Exception {
+ byte[] fileData = new byte[100_000_000];
+
+ runProcessor(fileData);
+
+ assertSuccess(DIRECTORY, FILE_NAME, fileData);
}
@Test
- public void testPutFile() throws Exception {
+ public void testPutFileWithNonExistingFileSystem() {
+ runner.setProperty(PutAzureDataLakeStorage.FILESYSTEM, "dummy");
+
+ runProcessor(FILE_DATA);
+
+ assertFailure();
+ }
+
+ @Test
+ public void testPutFileWithInvalidDirectory() {
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, "/dir1");
+
+ runProcessor(FILE_DATA);
+
+ assertFailure();
+ }
+
+ @Test
+ public void testPutFileWithInvalidFileName() {
+ runner.setProperty(PutAzureDataLakeStorage.FILE, "/file1");
+
+ runProcessor(FILE_DATA);
+
+ assertFailure();
+ }
+
+ @Test
+ public void testPutFileWithSpacesInDirectoryAndFileName() throws Exception
{
+ String directory = "dir 1";
+ String fileName = "file 1";
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY, directory);
+ runner.setProperty(PutAzureDataLakeStorage.FILE, fileName);
+
+ runProcessor(FILE_DATA);
+
+ assertSuccess(directory, fileName, FILE_DATA);
+ }
+
+ @Ignore
+ // the existing file gets overwritten without error
+ // seems to be a bug in the Azure lib
+ @Test
+ public void testPutFileToExistingFile() {
+ fileSystemClient.createFile(String.format("%s/%s", DIRECTORY,
FILE_NAME));
+
+ runProcessor(FILE_DATA);
+
+ assertFailure();
+ }
+
+ @Test
+ public void testPutFileWithEL() throws Exception {
+ Map<String, String> attributes = createAttributesMap();
+ setELProperties();
+
+ runProcessor(FILE_DATA, attributes);
+
+ assertSuccess(DIRECTORY, FILE_NAME, FILE_DATA);
+ }
+
+ @Test
+ public void testPutFileWithELButFilesystemIsNotSpecified() {
+ Map<String, String> attributes = createAttributesMap();
+ attributes.remove(EL_FILESYSTEM);
+ setELProperties();
+
+ runProcessor(FILE_DATA, attributes);
+
+ assertFailure();
+ }
+
+ @Test
+ public void testPutFileWithELButFileNameIsNotSpecified() {
+ Map<String, String> attributes = createAttributesMap();
+ attributes.remove(EL_FILE_NAME);
+ setELProperties();
+
+ runProcessor(FILE_DATA, attributes);
+
+ assertFailure();
+ }
+
+ private Map<String, String> createAttributesMap() {
+ Map<String, String> attributes = new HashMap<>();
+
+ attributes.put(EL_FILESYSTEM, fileSystemName);
+ attributes.put(EL_DIRECTORY, DIRECTORY);
+ attributes.put(EL_FILE_NAME, FILE_NAME);
+
+ return attributes;
+ }
+
+ private void setELProperties() {
+ runner.setProperty(PutAzureDataLakeStorage.FILESYSTEM,
String.format("${%s}", EL_FILESYSTEM));
+ runner.setProperty(PutAzureDataLakeStorage.DIRECTORY,
String.format("${%s}", EL_DIRECTORY));
+ runner.setProperty(PutAzureDataLakeStorage.FILE,
String.format("${%s}", EL_FILE_NAME));
+ }
+
+ private void runProcessor(byte[] fileData) {
+ runProcessor(fileData, Collections.emptyMap());
+ }
+
+ private void runProcessor(byte[] testData, Map<String, String> attributes)
{
runner.assertValid();
- runner.enqueue("0123456789".getBytes());
+ runner.enqueue(testData, attributes);
runner.run();
-
- assertResult();
}
+ private void assertSuccess(String directory, String fileName, byte[]
fileData) throws Exception {
+ assertFlowFile(directory, fileName, fileData);
+ assertAzureFile(directory, fileName, fileData);
+ }
- private void assertResult() throws Exception {
+ private void assertFlowFile(String directory, String fileName, byte[]
fileData) throws Exception {
runner.assertAllFlowFilesTransferred(PutAzureDataLakeStorage.REL_SUCCESS, 1);
- List<MockFlowFile> flowFilesForRelationship =
runner.getFlowFilesForRelationship(PutAzureDataLakeStorage.REL_SUCCESS);
- for (MockFlowFile flowFile : flowFilesForRelationship) {
- flowFile.assertContentEquals("0123456789".getBytes());
- flowFile.assertAttributeEquals("azure.length", "10");
+
+ MockFlowFile flowFile =
runner.getFlowFilesForRelationship(PutAzureDataLakeStorage.REL_SUCCESS).get(0);
+
+ flowFile.assertContentEquals(fileData);
+
+ flowFile.assertAttributeEquals("azure.filesystem", fileSystemName);
+ flowFile.assertAttributeEquals("azure.directory", directory);
+ flowFile.assertAttributeEquals("azure.filename", fileName);
+
+ String urlEscapedDirectory =
UrlEscapers.urlPathSegmentEscaper().escape(directory);
+ String urlEscapedFileName =
UrlEscapers.urlPathSegmentEscaper().escape(fileName);
+ String primaryUri = StringUtils.isNotEmpty(directory)
+ ? String.format("https://%s.dfs.core.windows.net/%s/%s/%s",
getAccountName(), fileSystemName, urlEscapedDirectory, urlEscapedFileName)
+ : String.format("https://%s.dfs.core.windows.net/%s/%s",
getAccountName(), fileSystemName, urlEscapedFileName);
+ flowFile.assertAttributeEquals("azure.primaryUri", primaryUri);
+
+ flowFile.assertAttributeEquals("azure.length",
Integer.toString(fileData.length));
+ }
+
+ private void assertAzureFile(String directory, String fileName, byte[]
fileData) {
+ DataLakeFileClient fileClient;
+ if (StringUtils.isNotEmpty(directory)) {
+ DataLakeDirectoryClient directoryClient =
fileSystemClient.getDirectoryClient(directory);
+ assertTrue(directoryClient.exists());
+
+ fileClient = directoryClient.getFileClient(fileName);
+ } else {
+ fileClient = fileSystemClient.getFileClient(fileName);
}
+ assertTrue(fileClient.exists());
+ assertEquals(fileData.length,
fileClient.getProperties().getFileSize());
+ }
+
+ private void assertFailure() {
+
runner.assertAllFlowFilesTransferred(PutAzureDataLakeStorage.REL_FAILURE, 1);
}
}
diff --git
a/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
new file mode 100644
index 0000000..59800bb
--- /dev/null
+++
b/nifi-nar-bundles/nifi-azure-bundle/nifi-azure-processors/src/test/java/org/apache/nifi/processors/azure/storage/TestAbstractAzureDataLakeStorage.java
@@ -0,0 +1,124 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.processors.azure.storage;
+
+import org.apache.nifi.util.TestRunner;
+import org.apache.nifi.util.TestRunners;
+import org.junit.Before;
+import org.junit.Test;
+
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.ACCOUNT_KEY;
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.ACCOUNT_NAME;
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.DIRECTORY;
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.FILE;
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.FILESYSTEM;
+import static
org.apache.nifi.processors.azure.AbstractAzureDataLakeStorageProcessor.SAS_TOKEN;
+
+public class TestAbstractAzureDataLakeStorage {
+
+ private TestRunner runner;
+
+ @Before
+ public void setUp() {
+ // test the property validation in the abstract class via the put
processor
+ runner = TestRunners.newTestRunner(PutAzureDataLakeStorage.class);
+
+ runner.setProperty(ACCOUNT_NAME, "accountName");
+ runner.setProperty(ACCOUNT_KEY, "accountKey");
+ runner.setProperty(FILESYSTEM, "filesystem");
+ runner.setProperty(DIRECTORY, "directory");
+ runner.setProperty(FILE, "file");
+ }
+
+ @Test
+ public void testValidWhenAccountNameAndAccountKeySpecified() {
+ runner.assertValid();
+ }
+
+ @Test
+ public void testValidWhenAccountNameAndSasTokenSpecified() {
+ runner.removeProperty(ACCOUNT_KEY);
+ runner.setProperty(SAS_TOKEN, "sasToken");
+
+ runner.assertValid();
+ }
+
+ @Test
+ public void testNotValidWhenNoAccountNameSpecified() {
+ runner.removeProperty(ACCOUNT_NAME);
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testNotValidWhenNoAccountKeyNorSasTokenSpecified() {
+ runner.removeProperty(ACCOUNT_KEY);
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testNotValidWhenBothAccountKeyAndSasTokenSpecified() {
+ runner.setProperty(SAS_TOKEN, "sasToken");
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testNotValidWhenNoFilesystemSpecified() {
+ runner.removeProperty(FILESYSTEM);
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testNotValidWhenFilesystemIsEmptyString() {
+ runner.setProperty(FILESYSTEM, "");
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testNotValidWhenNoDirectorySpecified() {
+ runner.removeProperty(DIRECTORY);
+
+ runner.assertNotValid();
+ }
+
+ @Test
+ public void testValidWhenDirectoryIsEmptyString() {
+ // the empty string is for the filesystem root directory
+ runner.setProperty(DIRECTORY, "");
+
+ runner.assertValid();
+ }
+
+ @Test
+ public void testValidWhenNoFileSpecified() {
+ // the default value will be used
+ runner.removeProperty(FILE);
+
+ runner.assertValid();
+ }
+
+ @Test
+ public void testNotValidWhenFileIsEmptyString() {
+ runner.setProperty(FILE, "");
+
+ runner.assertNotValid();
+ }
+}