This is an automated email from the ASF dual-hosted git repository.

lordgamez pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi-minifi-cpp.git

commit 9984cf9baa7d93468c7923e60bea5be613d0eb3f
Author: Adam Debreceni <[email protected]>
AuthorDate: Wed Nov 2 12:02:16 2022 +0100

    MINIFICPP-1976 - Do not store a whole Relationship object for each transfer
    
    Signed-off-by: Gabor Gyimesi <[email protected]>
    
    This closes #1447
---
 libminifi/include/core/ProcessSession.h | 22 ++++++--
 libminifi/include/core/ProcessorNode.h  |  9 +--
 libminifi/src/SchedulingAgent.cpp       | 10 +---
 libminifi/src/core/ProcessSession.cpp   | 97 +++++++++++++++++++--------------
 4 files changed, 76 insertions(+), 62 deletions(-)

diff --git a/libminifi/include/core/ProcessSession.h 
b/libminifi/include/core/ProcessSession.h
index 702e35231..208f7f100 100644
--- a/libminifi/include/core/ProcessSession.h
+++ b/libminifi/include/core/ProcessSession.h
@@ -28,6 +28,7 @@
 #include <algorithm>
 #include <set>
 #include <unordered_map>
+#include <unordered_set>
 
 #include "ProcessContext.h"
 #include "FlowFileRecord.h"
@@ -161,16 +162,25 @@ class ProcessSession : public ReferenceContainer {
     std::shared_ptr<FlowFile> snapshot;
   };
 
+  using Relationships = std::unordered_set<Relationship>;
+
+  Relationships relationships_;
+
+  struct NewFlowFileInfo {
+    std::shared_ptr<core::FlowFile> flow_file;
+    const Relationship* rel{nullptr};
+  };
+
   // FlowFiles being modified by current process session
-  std::map<utils::Identifier, FlowFileUpdate> _updatedFlowFiles;
+  std::map<utils::Identifier, FlowFileUpdate> updated_flowfiles_;
+  // updated FlowFiles being transferred to the relationship
+  std::map<utils::Identifier, const Relationship*> updated_relationships_;
   // FlowFiles being added by current process session
-  std::map<utils::Identifier, std::shared_ptr<core::FlowFile>> _addedFlowFiles;
+  std::map<utils::Identifier, NewFlowFileInfo> added_flowfiles_;
   // FlowFiles being deleted by current process session
-  std::vector<std::shared_ptr<core::FlowFile>> _deletedFlowFiles;
-  // FlowFiles being transfered to the relationship
-  std::map<utils::Identifier, Relationship> _transferRelationship;
+  std::vector<std::shared_ptr<core::FlowFile>> deleted_flowfiles_;
   // FlowFiles being cloned for multiple connections per relationship
-  std::vector<std::shared_ptr<core::FlowFile>> _clonedFlowFiles;
+  std::vector<std::shared_ptr<core::FlowFile>> cloned_flowfiles_;
 
  private:
   enum class RouteResult {
diff --git a/libminifi/include/core/ProcessorNode.h 
b/libminifi/include/core/ProcessorNode.h
index 7afce4ceb..96014ed9d 100644
--- a/libminifi/include/core/ProcessorNode.h
+++ b/libminifi/include/core/ProcessorNode.h
@@ -15,8 +15,7 @@
  * limitations under the License.
  */
 
-#ifndef LIBMINIFI_INCLUDE_CORE_PROCESSORNODE_H_
-#define LIBMINIFI_INCLUDE_CORE_PROCESSORNODE_H_
+#pragma once
 
 #include <memory>
 #include <set>
@@ -167,11 +166,11 @@ class ProcessorNode : public ConfigurableComponent, 
public Connectable {
     processor_->setAutoTerminatedRelationships(relationships);
   }
 
-  bool isAutoTerminated(Relationship relationship) {
+  bool isAutoTerminated(const Relationship& relationship) {
     return processor_->isAutoTerminated(relationship);
   }
 
-  bool isSupportedRelationship(Relationship relationship) {
+  bool isSupportedRelationship(const Relationship& relationship) {
     return processor_->isSupportedRelationship(relationship);
   }
 
@@ -278,5 +277,3 @@ class ProcessorNode : public ConfigurableComponent, public 
Connectable {
 };
 
 }  // namespace org::apache::nifi::minifi::core
-
-#endif  // LIBMINIFI_INCLUDE_CORE_PROCESSORNODE_H_
diff --git a/libminifi/src/SchedulingAgent.cpp 
b/libminifi/src/SchedulingAgent.cpp
index 5ad80fc9a..5b59faa31 100644
--- a/libminifi/src/SchedulingAgent.cpp
+++ b/libminifi/src/SchedulingAgent.cpp
@@ -32,10 +32,7 @@ bool hasWorkToDo(org::apache::nifi::minifi::core::Processor* 
processor) {
 }
 }  // namespace
 
-namespace org {
-namespace apache {
-namespace nifi {
-namespace minifi {
+namespace org::apache::nifi::minifi {
 
 std::future<utils::TaskRescheduleInfo> 
SchedulingAgent::enableControllerService(std::shared_ptr<core::controller::ControllerServiceNode>
 &serviceNode) {
   logger_->log_info("Enabling CSN in SchedulingAgent %s", 
serviceNode->getName());
@@ -141,7 +138,4 @@ void SchedulingAgent::watchDogFunc() {
   }
 }
 
-}  // namespace minifi
-}  // namespace nifi
-}  // namespace apache
-}  // namespace org
+}  // namespace org::apache::nifi::minifi
diff --git a/libminifi/src/core/ProcessSession.cpp 
b/libminifi/src/core/ProcessSession.cpp
index d472fff0e..0181f1e36 100644
--- a/libminifi/src/core/ProcessSession.cpp
+++ b/libminifi/src/core/ProcessSession.cpp
@@ -85,10 +85,10 @@ ProcessSession::~ProcessSession() {
 
 void ProcessSession::add(const std::shared_ptr<core::FlowFile> &record) {
   utils::Identifier uuid = record->getUUID();
-  if (_updatedFlowFiles.find(uuid) != _updatedFlowFiles.end()) {
+  if (updated_flowfiles_.find(uuid) != updated_flowfiles_.end()) {
     throw Exception(ExceptionType::PROCESSOR_EXCEPTION, "Mustn't add file that 
was provided by this session");
   }
-  _addedFlowFiles[uuid] = record;
+  added_flowfiles_[uuid].flow_file = record;
   record->setDeleted(false);
 }
 
@@ -114,7 +114,7 @@ std::shared_ptr<core::FlowFile> 
ProcessSession::create(const std::shared_ptr<cor
   }
 
   utils::Identifier uuid = record->getUUID();
-  _addedFlowFiles[uuid] = record;
+  added_flowfiles_[uuid].flow_file = record;
   logger_->log_debug("Create FlowFile with UUID %s", record->getUUIDStr());
   std::stringstream details;
   details << process_context_->getProcessorNode()->getName() << " creates flow 
record " << record->getUUIDStr();
@@ -146,7 +146,7 @@ std::shared_ptr<core::FlowFile> 
ProcessSession::cloneDuringTransfer(const std::s
   if (flow_version != nullptr) {
     record->setAttribute(SpecialFlowAttribute::FLOW_ID, 
flow_version->getFlowId());
   }
-  this->_clonedFlowFiles.push_back(record);
+  this->cloned_flowfiles_.push_back(record);
   logger_->log_debug("Clone FlowFile with UUID %s during transfer", 
record->getUUIDStr());
   // Copy attributes
   for (const auto& attribute : parent->getAttributes()) {
@@ -196,7 +196,7 @@ std::shared_ptr<core::FlowFile> ProcessSession::clone(const 
std::shared_ptr<core
 
 void ProcessSession::remove(const std::shared_ptr<core::FlowFile> &flow) {
   flow->setDeleted(true);
-  _deletedFlowFiles.push_back(flow);
+  deleted_flowfiles_.push_back(flow);
   std::string reason = process_context_->getProcessorNode()->getName() + " 
drop flow record " + flow->getUUIDStr();
   provenance_report_->drop(flow, reason);
 }
@@ -224,14 +224,18 @@ void ProcessSession::penalize(const 
std::shared_ptr<core::FlowFile> &flow) {
 void ProcessSession::transfer(const std::shared_ptr<core::FlowFile>& flow, 
const Relationship& relationship) {
   logging::LOG_INFO(logger_) << "Transferring " << flow->getUUIDStr() << " 
from " << process_context_->getProcessorNode()->getName() << " to relationship 
" << relationship.getName();
   utils::Identifier uuid = flow->getUUID();
-  _transferRelationship[uuid] = relationship;
+  if (auto it = added_flowfiles_.find(uuid); it != added_flowfiles_.end()) {
+    it->second.rel = &*relationships_.insert(relationship).first;
+  } else {
+    updated_relationships_[uuid] = &*relationships_.insert(relationship).first;
+  }
   flow->setDeleted(false);
 }
 
 void ProcessSession::write(const std::shared_ptr<core::FlowFile> &flow, const 
io::OutputStreamCallback& callback) {
-  gsl_ExpectsAudit(_updatedFlowFiles.contains(flow->getUUID())
-      || _addedFlowFiles.contains(flow->getUUID())
-      || std::any_of(_clonedFlowFiles.begin(), _clonedFlowFiles.end(), 
[&flow](const auto& flow_file) { return flow == flow_file; }));
+  gsl_ExpectsAudit(updated_flowfiles_.contains(flow->getUUID())
+      || added_flowfiles_.contains(flow->getUUID())
+      || std::any_of(cloned_flowfiles_.begin(), cloned_flowfiles_.end(), 
[&flow](const auto& flow_file) { return flow == flow_file; }));
 
   std::shared_ptr<ResourceClaim> claim = content_session_->create();
 
@@ -274,9 +278,9 @@ void ProcessSession::writeBuffer(const 
std::shared_ptr<core::FlowFile>& flow_fil
 }
 
 void ProcessSession::append(const std::shared_ptr<core::FlowFile> &flow, const 
io::OutputStreamCallback& callback) {
-  gsl_ExpectsAudit(_updatedFlowFiles.contains(flow->getUUID())
-      || _addedFlowFiles.contains(flow->getUUID())
-      || std::any_of(_clonedFlowFiles.begin(), _clonedFlowFiles.end(), 
[&flow](const auto& flow_file) { return flow == flow_file; }));
+  gsl_ExpectsAudit(updated_flowfiles_.contains(flow->getUUID())
+      || added_flowfiles_.contains(flow->getUUID())
+      || std::any_of(cloned_flowfiles_.begin(), cloned_flowfiles_.end(), 
[&flow](const auto& flow_file) { return flow == flow_file; }));
 
   std::shared_ptr<ResourceClaim> claim = flow->getResourceClaim();
   if (!claim) {
@@ -748,11 +752,15 @@ ProcessSession::RouteResult 
ProcessSession::routeFlowFile(const std::shared_ptr<
     return RouteResult::Ok_Deleted;
   }
   utils::Identifier uuid = record->getUUID();
-  auto itRelationship = _transferRelationship.find(uuid);
-  if (itRelationship == _transferRelationship.end()) {
+  Relationship relationship;
+  if (auto it = updated_relationships_.find(uuid); it != 
updated_relationships_.end()) {
+    gsl_Expects(it->second);
+    relationship = *it->second;
+  } else if (auto new_it = added_flowfiles_.find(uuid); new_it != 
added_flowfiles_.end() && new_it->second.rel) {
+    relationship = *new_it->second.rel;
+  } else {
     return RouteResult::Error_NoRelationship;
   }
-  Relationship relationship = itRelationship->second;
   // Find the relationship, we need to find the connections for that 
relationship
   const auto connections = 
process_context_->getProcessorNode()->getOutGoingConnections(relationship.getName());
   if (connections.empty()) {
@@ -798,7 +806,7 @@ void ProcessSession::commit() {
       transfers[relationship.getName()].transfer_size += record.getSize();
     };
     // First we clone the flow record based on the transferred relationship 
for updated flow record
-    for (auto && it : _updatedFlowFiles) {
+    for (auto && it : updated_flowfiles_) {
       auto record = it.second.modified;
       if (routeFlowFile(record, increaseTransferMetrics) == 
RouteResult::Error_NoRelationship) {
         // Can not find relationship for the flow
@@ -807,8 +815,8 @@ void ProcessSession::commit() {
     }
 
     // Do the same thing for added flow file
-    for (const auto& it : _addedFlowFiles) {
-      auto record = it.second;
+    for (const auto& it : added_flowfiles_) {
+      auto record = it.second.flow_file;
       if (routeFlowFile(record, increaseTransferMetrics) == 
RouteResult::Error_NoRelationship) {
         // Can not find relationship for the flow
         throw Exception(PROCESS_SESSION_EXCEPTION, "Can not find the transfer 
relationship for the added flow " + record->getUUIDStr());
@@ -819,9 +827,9 @@ void ProcessSession::commit() {
 
     Connectable* connection = nullptr;
     // Complete process the added and update flow files for the session, send 
the flow file to its queue
-    for (const auto &it : _updatedFlowFiles) {
+    for (const auto &it : updated_flowfiles_) {
       auto record = it.second.modified;
-      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"_updatedFlowFiles");
+      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"updated_flowfiles_");
       if (record->isDeleted()) {
         continue;
       }
@@ -831,9 +839,9 @@ void ProcessSession::commit() {
         connectionQueues[connection].push_back(record);
       }
     }
-    for (const auto &it : _addedFlowFiles) {
-      auto record = it.second;
-      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"_addedFlowFiles");
+    for (const auto &it : added_flowfiles_) {
+      auto record = it.second.flow_file;
+      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"added_flowfiles_");
       if (record->isDeleted()) {
         continue;
       }
@@ -843,8 +851,8 @@ void ProcessSession::commit() {
       }
     }
     // Process the clone flow files
-    for (const auto &record : _clonedFlowFiles) {
-      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"_clonedFlowFiles");
+    for (const auto &record : cloned_flowfiles_) {
+      logger_->log_trace("See %s in %s", record->getUUIDStr(), 
"cloned_flowfiles_");
       if (record->isDeleted()) {
         continue;
       }
@@ -854,7 +862,7 @@ void ProcessSession::commit() {
       }
     }
 
-    for (const auto& record : _deletedFlowFiles) {
+    for (const auto& record : deleted_flowfiles_) {
       if (!record->isDeleted()) {
         continue;
       }
@@ -872,7 +880,7 @@ void ProcessSession::commit() {
       throw Exception(PROCESS_SESSION_EXCEPTION, "State manager commit 
failed.");
     }
 
-    persistFlowFilesBeforeTransfer(connectionQueues, _updatedFlowFiles);
+    persistFlowFilesBeforeTransfer(connectionQueues, updated_flowfiles_);
 
     for (auto& cq : connectionQueues) {
       auto connection = dynamic_cast<Connection*>(cq.first);
@@ -894,12 +902,13 @@ void ProcessSession::commit() {
     }
 
     // All done
-    _updatedFlowFiles.clear();
-    _addedFlowFiles.clear();
-    _clonedFlowFiles.clear();
-    _deletedFlowFiles.clear();
+    updated_flowfiles_.clear();
+    added_flowfiles_.clear();
+    cloned_flowfiles_.clear();
+    deleted_flowfiles_.clear();
 
-    _transferRelationship.clear();
+    updated_relationships_.clear();
+    relationships_.clear();
     // persistent the provenance report
     this->provenance_report_->commit();
     logger_->log_trace("ProcessSession committed for %s", 
process_context_->getProcessorNode()->getName());
@@ -919,7 +928,7 @@ void ProcessSession::rollback() {
 
   try {
     // Requeue the snapshot of the flowfile back
-    for (const auto &it : _updatedFlowFiles) {
+    for (const auto &it : updated_flowfiles_) {
       auto flowFile = it.second.modified;
       // restore flowFile to original state
       *flowFile = *it.second.snapshot;
@@ -931,7 +940,7 @@ void ProcessSession::rollback() {
       connectionQueues[flowFile->getConnection()].push_back(flowFile);
     }
 
-    for (const auto& record : _deletedFlowFiles) {
+    for (const auto& record : deleted_flowfiles_) {
       record->setDeleted(false);
     }
 
@@ -953,10 +962,11 @@ void ProcessSession::rollback() {
       throw Exception(PROCESS_SESSION_EXCEPTION, "State manager rollback 
failed.");
     }
 
-    _clonedFlowFiles.clear();
-    _addedFlowFiles.clear();
-    _updatedFlowFiles.clear();
-    _deletedFlowFiles.clear();
+    cloned_flowfiles_.clear();
+    added_flowfiles_.clear();
+    updated_flowfiles_.clear();
+    deleted_flowfiles_.clear();
+    relationships_.clear();
     logger_->log_warn("ProcessSession rollback for %s executed", 
process_context_->getProcessorNode()->getName());
   } catch (const std::exception& exception) {
     logger_->log_warn("Caught Exception during process session rollback, type: 
%s, what: %s", typeid(exception).name(), exception.what());
@@ -1097,7 +1107,7 @@ std::shared_ptr<core::FlowFile> ProcessSession::get() {
       *snapshot = *ret;
       logger_->log_debug("Create Snapshot FlowFile with UUID %s", 
snapshot->getUUIDStr());
       utils::Identifier uuid = ret->getUUID();
-      _updatedFlowFiles[uuid] = {ret, snapshot};
+      updated_flowfiles_[uuid] = {ret, snapshot};
       auto flow_version = 
process_context_->getProcessorNode()->getFlowIdentifier();
       if (flow_version != nullptr) {
         ret->setAttribute(SpecialFlowAttribute::FLOW_ID, 
flow_version->getFlowId());
@@ -1127,9 +1137,12 @@ bool ProcessSession::outgoingConnectionsFull(const 
std::string& relationship) {
 }
 
 bool ProcessSession::existsFlowFileInRelationship(const Relationship 
&relationship) {
-  return std::any_of(_transferRelationship.begin(), 
_transferRelationship.end(),
-      [&relationship](const std::map<utils::Identifier, 
Relationship>::value_type &key_value_pair) {
-        return relationship == key_value_pair.second;
+  return std::any_of(updated_relationships_.begin(), 
updated_relationships_.end(),
+      [&](const auto& key_value_pair) {
+        return key_value_pair.second && relationship == *key_value_pair.second;
+  }) || std::any_of(added_flowfiles_.begin(), added_flowfiles_.end(),
+      [&](const auto& key_value_pair) {
+        return key_value_pair.second.rel && relationship == 
*key_value_pair.second.rel;
   });
 }
 

Reply via email to