Repository: nifi Updated Branches: refs/heads/master e4e263c29 -> 72ccc252f
NIFI-927: Use a serializable version of NiFiDataPacket in the spark receiver Project: http://git-wip-us.apache.org/repos/asf/nifi/repo Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/72ccc252 Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/72ccc252 Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/72ccc252 Branch: refs/heads/master Commit: 72ccc252fe118c51ccd3ddb6cba7df6b3d85397a Parents: e4e263c Author: Mark Payne <marka...@hotmail.com> Authored: Fri Sep 4 12:22:19 2015 -0400 Committer: Mark Payne <marka...@hotmail.com> Committed: Fri Sep 4 12:22:19 2015 -0400 ---------------------------------------------------------------------- .../org/apache/nifi/spark/NiFiReceiver.java | 13 +----- .../nifi/spark/StandardNiFiDataPacket.java | 43 ++++++++++++++++++++ 2 files changed, 44 insertions(+), 12 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/nifi/blob/72ccc252/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/NiFiReceiver.java ---------------------------------------------------------------------- diff --git a/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/NiFiReceiver.java b/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/NiFiReceiver.java index 8cbf60c..689abd0 100644 --- a/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/NiFiReceiver.java +++ b/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/NiFiReceiver.java @@ -170,18 +170,7 @@ public class NiFiReceiver extends Receiver<NiFiDataPacket> { StreamUtils.fillBuffer(inStream, data); final Map<String, String> attributes = dataPacket.getAttributes(); - final NiFiDataPacket NiFiDataPacket = new NiFiDataPacket() { - @Override - public byte[] getContent() { - return data; - } - - @Override - public Map<String, String> getAttributes() { - return attributes; - } - }; - + final NiFiDataPacket NiFiDataPacket = new StandardNiFiDataPacket(data, attributes); dataPackets.add(NiFiDataPacket); dataPacket = transaction.receive(); } while (dataPacket != null); http://git-wip-us.apache.org/repos/asf/nifi/blob/72ccc252/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/StandardNiFiDataPacket.java ---------------------------------------------------------------------- diff --git a/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/StandardNiFiDataPacket.java b/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/StandardNiFiDataPacket.java new file mode 100644 index 0000000..8f5e0bc --- /dev/null +++ b/nifi-external/nifi-spark-receiver/src/main/java/org/apache/nifi/spark/StandardNiFiDataPacket.java @@ -0,0 +1,43 @@ +/* + * 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.spark; + +import java.io.Serializable; +import java.util.Map; + +public class StandardNiFiDataPacket implements NiFiDataPacket, Serializable { + private static final long serialVersionUID = 6364005260220243322L; + + private final byte[] content; + private final Map<String, String> attributes; + + public StandardNiFiDataPacket(final byte[] content, final Map<String, String> attributes) { + this.content = content; + this.attributes = attributes; + } + + @Override + public byte[] getContent() { + return content; + } + + @Override + public Map<String, String> getAttributes() { + return attributes; + } + +}