This is an automated email from the ASF dual-hosted git repository.
joewitt pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git
The following commit(s) were added to refs/heads/main by this push:
new 27e7d0b NIFI-7707 - This closes #4451. fixed Publish GCP PubSub
processor for binary data
27e7d0b is described below
commit 27e7d0b02ab353989857662a5e9a94b11ed7691c
Author: Pierre Villard <[email protected]>
AuthorDate: Wed Aug 5 12:21:12 2020 +0200
NIFI-7707 - This closes #4451. fixed Publish GCP PubSub processor for
binary data
Signed-off-by: Joe Witt <[email protected]>
---
.../java/org/apache/nifi/processors/gcp/pubsub/PublishGCPubSub.java | 6 +++++-
1 file changed, 5 insertions(+), 1 deletion(-)
diff --git
a/nifi-nar-bundles/nifi-gcp-bundle/nifi-gcp-processors/src/main/java/org/apache/nifi/processors/gcp/pubsub/PublishGCPubSub.java
b/nifi-nar-bundles/nifi-gcp-bundle/nifi-gcp-processors/src/main/java/org/apache/nifi/processors/gcp/pubsub/PublishGCPubSub.java
index a41a610..5ad4e76 100644
---
a/nifi-nar-bundles/nifi-gcp-bundle/nifi-gcp-processors/src/main/java/org/apache/nifi/processors/gcp/pubsub/PublishGCPubSub.java
+++
b/nifi-nar-bundles/nifi-gcp-bundle/nifi-gcp-processors/src/main/java/org/apache/nifi/processors/gcp/pubsub/PublishGCPubSub.java
@@ -28,6 +28,8 @@ import com.google.pubsub.v1.ProjectTopicName;
import com.google.pubsub.v1.PubsubMessage;
import org.apache.nifi.annotation.behavior.DynamicProperty;
import org.apache.nifi.annotation.behavior.InputRequirement;
+import org.apache.nifi.annotation.behavior.SystemResource;
+import org.apache.nifi.annotation.behavior.SystemResourceConsideration;
import org.apache.nifi.annotation.behavior.InputRequirement.Requirement;
import org.apache.nifi.annotation.behavior.WritesAttribute;
import org.apache.nifi.annotation.behavior.WritesAttributes;
@@ -75,6 +77,8 @@ import static
org.apache.nifi.processors.gcp.pubsub.PubSubAttributes.TOPIC_NAME_
@WritesAttribute(attribute = MESSAGE_ID_ATTRIBUTE, description =
MESSAGE_ID_DESCRIPTION),
@WritesAttribute(attribute = TOPIC_NAME_ATTRIBUTE, description =
TOPIC_NAME_DESCRIPTION)
})
+@SystemResourceConsideration(resource = SystemResource.MEMORY, description =
"The entirety of the FlowFile's content "
+ + "will be read into memory to be sent as a PubSub message.")
public class PublishGCPubSub extends AbstractGCPubSubProcessor{
public static final PropertyDescriptor TOPIC_NAME = new
PropertyDescriptor.Builder()
@@ -154,7 +158,7 @@ public class PublishGCPubSub extends
AbstractGCPubSubProcessor{
try {
final ByteArrayOutputStream baos = new
ByteArrayOutputStream();
session.exportTo(flowFile, baos);
- final ByteString flowFileContent =
ByteString.copyFromUtf8(baos.toString());
+ final ByteString flowFileContent =
ByteString.copyFrom(baos.toByteArray());
PubsubMessage message =
PubsubMessage.newBuilder().setData(flowFileContent)
.setPublishTime(Timestamp.newBuilder().build())