From e065b81116a807616c7ed5bcec29e036b6d39903 Mon Sep 17 00:00:00 2001 From: Ivan Chupin Date: Mon, 30 Sep 2024 12:06:11 +0200 Subject: [PATCH 1/2] PERSISTENCE_MANAGER_INSTANCE unique among all the bolts tusks in NBWorker topology. --- .../nbworker-topology/nbworker-topology.tmpl | 4 ++- .../java/org/openkilda/wfm/AbstractBolt.java | 2 +- .../share/history/service/HistoryService.java | 2 ++ .../persistence/PersistenceManager.java | 14 +++++++++ .../persistence/tx/TransactionManager.java | 2 ++ .../HibernateSessionFactorySupplier.java | 6 ++-- .../HibernateGenericRepository.java | 8 +++++ .../HibernateHistoryFlowEventRepository.java | 23 ++++++++++++-- .../nbworker/bolts/HistoryOperationsBolt.java | 3 +- .../bolts/PersistenceOperationsBolt.java | 30 +++++++++++++++++-- 10 files changed, 83 insertions(+), 11 deletions(-) diff --git a/confd/templates/nbworker-topology/nbworker-topology.tmpl b/confd/templates/nbworker-topology/nbworker-topology.tmpl index 00d2ccb3a6d..9715183fabe 100644 --- a/confd/templates/nbworker-topology/nbworker-topology.tmpl +++ b/confd/templates/nbworker-topology/nbworker-topology.tmpl @@ -4,7 +4,7 @@ # topology configuration config: topology.parallelism: {{ getv "/kilda_storm_nb_worker_parallelism" }} - topology.workers: {{ getv "/kilda_storm_nbworker_workers_count" }} + topology.workers: 1 topology.spouts.parallelism: {{ getv "/kilda_storm_spout_parallelism" }} # spout definitions @@ -24,3 +24,5 @@ bolts: parallelism: 1 - id: "flows-operations-bolt" parallelism: {{ getv "/kilda_storm_nb_worker_flow_operations_parallelism" }} + - id: "history-operations-bolt" + parallelism: 3 \ No newline at end of file diff --git a/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/AbstractBolt.java b/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/AbstractBolt.java index fb02cde4a1d..5bf86c97ab7 100644 --- a/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/AbstractBolt.java +++ b/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/AbstractBolt.java @@ -55,7 +55,7 @@ public abstract class AbstractBolt extends BaseRichBolt { protected boolean active = false; @Getter - protected final PersistenceManager persistenceManager; + protected PersistenceManager persistenceManager; @Getter(AccessLevel.PROTECTED) private transient OutputCollector output; diff --git a/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/share/history/service/HistoryService.java b/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/share/history/service/HistoryService.java index 7039628a32e..90bfe49e295 100644 --- a/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/share/history/service/HistoryService.java +++ b/src-java/base-topology/base-storm-topology/src/main/java/org/openkilda/wfm/share/history/service/HistoryService.java @@ -144,6 +144,8 @@ public void store(PortEventData data) { */ public List listFlowEvents(String flowId, Instant timeFrom, Instant timeTo, int maxCount) { List result = new ArrayList<>(); + log.info("CHUPIN HistoryService, implementation in transactionManager: {}", + transactionManager.getImplementation()); transactionManager.doInTransaction(() -> flowEventRepository .findByFlowIdAndTimeFrame(flowId, timeFrom, timeTo, maxCount) .forEach(entry -> { diff --git a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/PersistenceManager.java b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/PersistenceManager.java index c19b1df0470..4b17d0a19c8 100644 --- a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/PersistenceManager.java +++ b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/PersistenceManager.java @@ -139,6 +139,20 @@ private PersistenceImplementation getImplementation(PersistenceImplementationTyp return implementation; } + /** + * Get retries limit. + */ + public int getTransactionRetriesLimit() { + return persistenceConfig.getTransactionRetriesLimit(); + } + + /** + * Get retries max delay. + */ + public int getTransactionRetriesMaxDelay() { + return persistenceConfig.getTransactionRetriesMaxDelay(); + } + public void install() { PersistenceContextManager.install(this); } diff --git a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/tx/TransactionManager.java b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/tx/TransactionManager.java index f3e1463608a..f0b2eec0e42 100644 --- a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/tx/TransactionManager.java +++ b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/tx/TransactionManager.java @@ -22,6 +22,7 @@ import org.openkilda.persistence.exceptions.PersistenceException; import org.openkilda.persistence.exceptions.RecoverablePersistenceException; +import lombok.Getter; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import net.jodah.failsafe.Failsafe; @@ -37,6 +38,7 @@ */ @Slf4j public class TransactionManager implements Serializable { + @Getter private final PersistenceImplementation implementation; private final int transactionRetriesLimit; diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/HibernateSessionFactorySupplier.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/HibernateSessionFactorySupplier.java index d7dccb196be..7f96fa66243 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/HibernateSessionFactorySupplier.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/HibernateSessionFactorySupplier.java @@ -34,7 +34,7 @@ import org.hibernate.boot.registry.StandardServiceRegistryBuilder; import org.hibernate.cfg.AvailableSettings; import org.hibernate.context.internal.ManagedSessionContext; -import org.hibernate.dialect.MySQLDialect; +import org.hibernate.dialect.MySQL8Dialect; import java.io.Serializable; import java.util.function.Supplier; @@ -62,14 +62,16 @@ public SessionFactory get() { } private SessionFactory makeHibernateSessionFactory() { + log.info("HibernateSessionFactorySupplier makeHibernateSessionFactory"); StandardServiceRegistry standardRegistry = new StandardServiceRegistryBuilder() .applySetting(AvailableSettings.DRIVER, hibernateConfig.getDriverClass()) .applySetting(AvailableSettings.USER, hibernateConfig.getUser()) .applySetting(AvailableSettings.PASS, hibernateConfig.getPassword()) .applySetting(AvailableSettings.URL, hibernateConfig.getUrl()) - .applySetting(AvailableSettings.DIALECT, MySQLDialect.class.getName()) + .applySetting(AvailableSettings.DIALECT, MySQL8Dialect.class.getName()) .applySetting(AvailableSettings.CURRENT_SESSION_CONTEXT_CLASS, ManagedSessionContext.class.getName()) .applySetting(AvailableSettings.C3P0_IDLE_TEST_PERIOD, 600) // seconds? + .applySetting(AvailableSettings.C3P0_MAX_SIZE, 30) .applySetting(AvailableSettings.C3P0_CONFIG_PREFIX + ".testConnectionOnCheckout", true) .applySetting(AvailableSettings.C3P0_CONFIG_PREFIX + ".preferredTestQuery", "SELECT 1") // TODO(surabujin): detect debugging mode and enable for it diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java index acdd99161ba..c77b28e7990 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java @@ -25,9 +25,11 @@ import org.openkilda.persistence.repositories.Repository; import org.openkilda.persistence.tx.TransactionManager; +import lombok.extern.slf4j.Slf4j; import org.hibernate.Session; import org.hibernate.SessionFactory; +@Slf4j public abstract class HibernateGenericRepository, V, H extends V> implements Repository { protected final HibernatePersistenceImplementation implementation; @@ -90,6 +92,12 @@ protected TransactionManager getTransactionManager() { return manager.getTransactionManager(implementation.getType()); } + protected TransactionManager getTransactionManager1() { + PersistenceManager pm = PersistenceContextManager.INSTANCE.getPersistenceManager(); + return new TransactionManager(implementation, pm.getTransactionRetriesLimit(), + pm.getTransactionRetriesMaxDelay()); + } + protected HibernateContextExtension getContextExtension() { PersistenceContext context = PersistenceContextManager.INSTANCE.getContextCreateIfMissing(); return implementation.getContextExtension(context); diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java index 07a55affaec..811b0abc214 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java @@ -31,6 +31,9 @@ import org.openkilda.persistence.hibernate.entities.history.HibernateFlowEvent_; import org.openkilda.persistence.hibernate.utils.UniqueKeyUtil; import org.openkilda.persistence.repositories.history.FlowEventRepository; +import org.openkilda.persistence.tx.TransactionManager; + +import lombok.extern.slf4j.Slf4j; import java.time.Instant; import java.util.ArrayList; @@ -43,6 +46,7 @@ import javax.persistence.criteria.Predicate; import javax.persistence.criteria.Root; +@Slf4j public class HibernateHistoryFlowEventRepository extends HibernateGenericRepository implements FlowEventRepository { @@ -60,10 +64,25 @@ public Optional findByTaskId(String taskId) { return getTransactionManager().doInTransaction(() -> findEntityByTaskId(taskId).map(FlowEvent::new)); } + /** + * Retrieves a list of {@link FlowEvent} objects filtered by the given flow ID and time frame. + * The method performs a transactional operation to fetch the events, maps them to new {@link FlowEvent} objects, + * and then reverses the order of the list before returning it. + * + * @param flowId The ID of the flow for which events are being retrieved. Cannot be null. + * @param timeFrom The start of the time frame for the events. + * @param timeTo The end of the time frame for the events. + * @param maxCount The maximum number of events to retrieve. + * @return A list of {@link FlowEvent} objects matching the given criteria, in reverse chronological order. + * If no events match the criteria, an empty list is returned. + */ @Override public List findByFlowIdAndTimeFrame( String flowId, Instant timeFrom, Instant timeTo, int maxCount) { - List results = getTransactionManager().doInTransaction( + TransactionManager transactionManager = getTransactionManager1(); + log.info("CHUPIN HibernateHistoryFlowEventRepository findByFlowIdAndTimeFrame, implementation: {}", + implementation); + List results = transactionManager.doInTransaction( () -> fetch(flowId, timeFrom, timeTo, maxCount).stream() .map(FlowEvent::new) .collect(Collectors.toList())); @@ -74,7 +93,7 @@ public List findByFlowIdAndTimeFrame( @Override public List findFlowStatusesByFlowIdAndTimeFrame( String flowId, Instant timeFrom, Instant timeTo, int maxCount) { - List results = getTransactionManager().doInTransaction( + List results = getTransactionManager1().doInTransaction( () -> fetch(flowId, timeFrom, timeTo, maxCount).stream() .flatMap(entry -> entry.getActions().stream()) .map(this::extractStatusUpdates) diff --git a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/HistoryOperationsBolt.java b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/HistoryOperationsBolt.java index aac07219d4e..c59d3d85335 100644 --- a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/HistoryOperationsBolt.java +++ b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/HistoryOperationsBolt.java @@ -57,8 +57,7 @@ public HistoryOperationsBolt(PersistenceManager persistenceManager) { @Override public void init() { super.init(); - - historyService = new HistoryService(persistenceManager); + historyService = new HistoryService(PERSISTENCE_MANAGER_INSTANCE); } @Override diff --git a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java index 8302310a91c..246970603be 100644 --- a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java +++ b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java @@ -27,18 +27,22 @@ import org.openkilda.wfm.AbstractBolt; import org.openkilda.wfm.topology.nbworker.StreamType; +import org.apache.storm.task.OutputCollector; +import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import java.util.List; +import java.util.Map; public abstract class PersistenceOperationsBolt extends AbstractBolt { public static final String FIELD_ID_REQUEST = "request"; protected transient RepositoryFactory repositoryFactory; protected transient TransactionManager transactionManager; + protected static volatile PersistenceManager PERSISTENCE_MANAGER_INSTANCE; PersistenceOperationsBolt(PersistenceManager persistenceManager) { super(persistenceManager); @@ -48,12 +52,32 @@ protected String getCorrelationId() { return getCommandContext().getCorrelationId(); } + /** + * Sample doc. + * + * @param stormConf The Storm configuration for this bolt. + * @param context This object can be used to get information about this task's place within the topology. + * @param collector The collector is used to emit tuples from this bolt. + */ + @Override + public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + synchronized (HistoryOperationsBolt.class) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + PERSISTENCE_MANAGER_INSTANCE = persistenceManager; + PERSISTENCE_MANAGER_INSTANCE.install(); + } + } + } + persistenceManager = null; + super.prepare(stormConf, context, collector); + } + @Override public void init() { super.init(); - - repositoryFactory = persistenceManager.getRepositoryFactory(); - transactionManager = persistenceManager.getTransactionManager(); + repositoryFactory = PERSISTENCE_MANAGER_INSTANCE.getRepositoryFactory(); + transactionManager = PERSISTENCE_MANAGER_INSTANCE.getTransactionManager(); } protected void handleInput(Tuple input) throws Exception { From 780746b96a1481dc5c1b523b6b56375b048df71f Mon Sep 17 00:00:00 2001 From: Ivan Chupin Date: Mon, 30 Sep 2024 13:57:48 +0200 Subject: [PATCH 2/2] PERSISTENCE_MANAGER_INSTANCE unique among all the bolts tusks in the following topologies: NBWorker topology, network topology, history topology. + removed unused methods in FlowEventRepository, FermaFlowEventRepository kilda_storm_nbworker_workers_count 2 -> 1 kilda_storm_network_workers_count 2->1 kilda_storm_history_workers_count 2->1 --- .../nbworker-topology/nbworker-topology.tmpl | 4 +-- confd/vars/main.yaml | 6 ++-- .../topology/history/bolts/HistoryBolt.java | 31 ++++++++++++++++++- .../history/FlowEventRepository.java | 3 -- .../HibernateGenericRepository.java | 7 +++-- .../HibernateHistoryFlowEventRepository.java | 14 +++------ ...HibernateHistoryHaFlowEventRepository.java | 2 +- .../HibernateHistoryPortEventRepository.java | 2 +- .../FermaFlowEventRepository.java | 11 ------- .../bolts/PersistenceOperationsBolt.java | 4 ++- .../storm/bolt/history/HistoryHandler.java | 30 +++++++++++++++++- 11 files changed, 78 insertions(+), 36 deletions(-) diff --git a/confd/templates/nbworker-topology/nbworker-topology.tmpl b/confd/templates/nbworker-topology/nbworker-topology.tmpl index 9715183fabe..2da8e8b2f87 100644 --- a/confd/templates/nbworker-topology/nbworker-topology.tmpl +++ b/confd/templates/nbworker-topology/nbworker-topology.tmpl @@ -4,7 +4,7 @@ # topology configuration config: topology.parallelism: {{ getv "/kilda_storm_nb_worker_parallelism" }} - topology.workers: 1 + topology.workers: {{ getv "/kilda_storm_nbworker_workers_count" }} topology.spouts.parallelism: {{ getv "/kilda_storm_spout_parallelism" }} # spout definitions @@ -25,4 +25,4 @@ bolts: - id: "flows-operations-bolt" parallelism: {{ getv "/kilda_storm_nb_worker_flow_operations_parallelism" }} - id: "history-operations-bolt" - parallelism: 3 \ No newline at end of file + parallelism: 3 diff --git a/confd/vars/main.yaml b/confd/vars/main.yaml index e6572968ca2..ec2d7a5d72c 100644 --- a/confd/vars/main.yaml +++ b/confd/vars/main.yaml @@ -194,7 +194,7 @@ kilda_storm_stats_workers_count: 2 # Network kilda_storm_network_parallelism: 2 -kilda_storm_network_workers_count: 2 +kilda_storm_network_workers_count: 1 # Reroute kilda_storm_reroute_parallelism: 2 @@ -207,7 +207,7 @@ kilda_storm_floodlight_router_workers_count: 2 # NB worker kilda_storm_nb_worker_parallelism: 2 kilda_storm_nb_worker_flow_operations_parallelism: 2 -kilda_storm_nbworker_workers_count: 2 +kilda_storm_nbworker_workers_count: 1 # Server42 control kilda_storm_server42_control_parallelism: 2 @@ -217,7 +217,7 @@ kilda_storm_server42_control_count: 2 kilda_storm_history_parallelism: 2 kilda_storm_history_bolt_parallelism: 2 kilda_storm_history_bolt_num_tasks: 2 -kilda_storm_history_workers_count: 2 +kilda_storm_history_workers_count: 1 # Connected devices kilda_storm_connected_devices_parallelism: 2 diff --git a/src-java/history-topology/history-storm-topology/src/main/java/org/openkilda/wfm/topology/history/bolts/HistoryBolt.java b/src-java/history-topology/history-storm-topology/src/main/java/org/openkilda/wfm/topology/history/bolts/HistoryBolt.java index 1a77722e25f..20f15dd836b 100644 --- a/src-java/history-topology/history-storm-topology/src/main/java/org/openkilda/wfm/topology/history/bolts/HistoryBolt.java +++ b/src-java/history-topology/history-storm-topology/src/main/java/org/openkilda/wfm/topology/history/bolts/HistoryBolt.java @@ -27,20 +27,49 @@ import org.openkilda.wfm.share.zk.ZkStreams; import org.openkilda.wfm.share.zk.ZooKeeperBolt; +import org.apache.storm.task.OutputCollector; +import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; +import java.util.Map; + public class HistoryBolt extends AbstractBolt { private transient HistoryService historyService; + private static volatile PersistenceManager PERSISTENCE_MANAGER_INSTANCE; public HistoryBolt(PersistenceManager persistenceManager, String lifeCycleEventSourceComponent) { super(persistenceManager, lifeCycleEventSourceComponent); } + + /** + * Called when a task for this component is initialized within a worker on the cluster. + * It provides the bolt with the environment in which the bolt executes. + * Static instance of PersistenceManager is initialized at this step. + * + * @param stormConf The Storm configuration for this bolt. + * @param context This object can be used to get information about this task's place within the topology. + * @param collector The collector is used to emit tuples from this bolt. + */ + @Override + public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + synchronized (HistoryBolt.class) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + PERSISTENCE_MANAGER_INSTANCE = persistenceManager; + PERSISTENCE_MANAGER_INSTANCE.install(); + } + } + } + persistenceManager = null; + super.prepare(stormConf, context, collector); + } + @Override protected void init() { - historyService = new HistoryService(persistenceManager); + historyService = new HistoryService(PERSISTENCE_MANAGER_INSTANCE); } @Override diff --git a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/repositories/history/FlowEventRepository.java b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/repositories/history/FlowEventRepository.java index c3f3f741bae..deb2d32d5cb 100644 --- a/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/repositories/history/FlowEventRepository.java +++ b/src-java/kilda-persistence-api/src/main/java/org/openkilda/persistence/repositories/history/FlowEventRepository.java @@ -21,13 +21,10 @@ import java.time.Instant; import java.util.List; -import java.util.Optional; public interface FlowEventRepository extends Repository { boolean existsByTaskId(String taskId); - Optional findByTaskId(String taskId); - List findByFlowIdAndTimeFrame(String flowId, Instant timeFrom, Instant timeTo, int maxCount); List findFlowStatusesByFlowIdAndTimeFrame(String flowId, Instant timeFrom, diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java index c77b28e7990..5e16bcf11f5 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateGenericRepository.java @@ -44,7 +44,7 @@ public void add(M model) { if (view instanceof EntityBase) { throw new IllegalArgumentException("Entity of class " + model + " already persisted"); } - getTransactionManager().doInTransaction(() -> { + getTransactionManagerWithLocalPersistenceImplementation().doInTransaction(() -> { H entity = makeEntity(view); getSession().persist(entity); model.setData(entity); @@ -65,7 +65,8 @@ public void remove(M model) { H hibernateView = (H) view; V detachedView = doDetach(model, hibernateView); - getTransactionManager().doInTransaction(() -> getSession().remove(hibernateView)); + getTransactionManagerWithLocalPersistenceImplementation() + .doInTransaction(() -> getSession().remove(hibernateView)); model.setData(detachedView); } @@ -92,7 +93,7 @@ protected TransactionManager getTransactionManager() { return manager.getTransactionManager(implementation.getType()); } - protected TransactionManager getTransactionManager1() { + protected TransactionManager getTransactionManagerWithLocalPersistenceImplementation() { PersistenceManager pm = PersistenceContextManager.INSTANCE.getPersistenceManager(); return new TransactionManager(implementation, pm.getTransactionRetriesLimit(), pm.getTransactionRetriesMaxDelay()); diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java index 811b0abc214..8954605709f 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryFlowEventRepository.java @@ -56,12 +56,8 @@ public HibernateHistoryFlowEventRepository(HibernatePersistenceImplementation im @Override public boolean existsByTaskId(String taskId) { - return getTransactionManager().doInTransaction(() -> findEntityByTaskId(taskId).isPresent()); - } - - @Override - public Optional findByTaskId(String taskId) { - return getTransactionManager().doInTransaction(() -> findEntityByTaskId(taskId).map(FlowEvent::new)); + return getTransactionManagerWithLocalPersistenceImplementation() + .doInTransaction(() -> findEntityByTaskId(taskId).isPresent()); } /** @@ -74,12 +70,12 @@ public Optional findByTaskId(String taskId) { * @param timeTo The end of the time frame for the events. * @param maxCount The maximum number of events to retrieve. * @return A list of {@link FlowEvent} objects matching the given criteria, in reverse chronological order. - * If no events match the criteria, an empty list is returned. + * If no events match the criteria, an empty list is returned. */ @Override public List findByFlowIdAndTimeFrame( String flowId, Instant timeFrom, Instant timeTo, int maxCount) { - TransactionManager transactionManager = getTransactionManager1(); + TransactionManager transactionManager = getTransactionManagerWithLocalPersistenceImplementation(); log.info("CHUPIN HibernateHistoryFlowEventRepository findByFlowIdAndTimeFrame, implementation: {}", implementation); List results = transactionManager.doInTransaction( @@ -93,7 +89,7 @@ public List findByFlowIdAndTimeFrame( @Override public List findFlowStatusesByFlowIdAndTimeFrame( String flowId, Instant timeFrom, Instant timeTo, int maxCount) { - List results = getTransactionManager1().doInTransaction( + List results = getTransactionManagerWithLocalPersistenceImplementation().doInTransaction( () -> fetch(flowId, timeFrom, timeTo, maxCount).stream() .flatMap(entry -> entry.getActions().stream()) .map(this::extractStatusUpdates) diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryHaFlowEventRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryHaFlowEventRepository.java index 70a8cbb44fa..fb5beb8508a 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryHaFlowEventRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryHaFlowEventRepository.java @@ -58,7 +58,7 @@ public Optional findEntityByTaskId(String taskId) { @Override public List findByHaFlowIdAndTimeFrame(String haFlowId, Instant timeFrom, Instant timeTo, int maxCount) { - List results = getTransactionManager().doInTransaction( + List results = getTransactionManagerWithLocalPersistenceImplementation().doInTransaction( () -> fetch(haFlowId, timeFrom, timeTo, maxCount).stream() .map(HaFlowEvent::new) .collect(Collectors.toList())); diff --git a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryPortEventRepository.java b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryPortEventRepository.java index 504d2ad0d63..89fd27a9b12 100644 --- a/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryPortEventRepository.java +++ b/src-java/kilda-persistence-hibernate/src/main/java/org/openkilda/persistence/hibernate/repositories/HibernateHistoryPortEventRepository.java @@ -42,7 +42,7 @@ public HibernateHistoryPortEventRepository(HibernatePersistenceImplementation im @Override public List findBySwitchIdAndPortNumber( SwitchId switchId, int portNumber, Instant start, Instant end) { - return getTransactionManager().doInTransaction( + return getTransactionManagerWithLocalPersistenceImplementation().doInTransaction( () -> findEntityBySwitchIdAndPortNumber(switchId, portNumber, start, end).stream() .map(PortEvent::new) .collect(Collectors.toList())); diff --git a/src-java/kilda-persistence-tinkerpop/src/main/java/org/openkilda/persistence/ferma/repositories/FermaFlowEventRepository.java b/src-java/kilda-persistence-tinkerpop/src/main/java/org/openkilda/persistence/ferma/repositories/FermaFlowEventRepository.java index 113d345c4f5..69323d3f7d4 100644 --- a/src-java/kilda-persistence-tinkerpop/src/main/java/org/openkilda/persistence/ferma/repositories/FermaFlowEventRepository.java +++ b/src-java/kilda-persistence-tinkerpop/src/main/java/org/openkilda/persistence/ferma/repositories/FermaFlowEventRepository.java @@ -35,7 +35,6 @@ import java.util.ArrayList; import java.util.Comparator; import java.util.List; -import java.util.Optional; import java.util.stream.Collectors; /** @@ -59,16 +58,6 @@ public boolean existsByTaskId(String taskId) { } } - @Override - public Optional findByTaskId(String taskId) { - List flowEventFrames = framedGraph().traverse(g -> g.V() - .hasLabel(FlowEventFrame.FRAME_LABEL) - .has(FlowEventFrame.TASK_ID_PROPERTY, taskId)) - .toListExplicit(FlowEventFrame.class); - return flowEventFrames.isEmpty() ? Optional.empty() : Optional.of(flowEventFrames.get(0)) - .map(FlowEvent::new); - } - @Override public List findByFlowIdAndTimeFrame(String flowId, Instant timeFrom, Instant timeTo, int maxCount) { return framedGraph().traverse(g -> { diff --git a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java index 246970603be..804aa532cf4 100644 --- a/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java +++ b/src-java/nbworker-topology/nbworker-storm-topology/src/main/java/org/openkilda/wfm/topology/nbworker/bolts/PersistenceOperationsBolt.java @@ -53,7 +53,9 @@ protected String getCorrelationId() { } /** - * Sample doc. + * Called when a task for this component is initialized within a worker on the cluster. + * It provides the bolt with the environment in which the bolt executes. + * Static instance of PersistenceManager is initialized at this step. * * @param stormConf The Storm configuration for this bolt. * @param context This object can be used to get information about this task's place within the topology. diff --git a/src-java/network-topology/network-storm-topology/src/main/java/org/openkilda/wfm/topology/network/storm/bolt/history/HistoryHandler.java b/src-java/network-topology/network-storm-topology/src/main/java/org/openkilda/wfm/topology/network/storm/bolt/history/HistoryHandler.java index 6431fb0c8f6..5d8bcf3dea1 100644 --- a/src-java/network-topology/network-storm-topology/src/main/java/org/openkilda/wfm/topology/network/storm/bolt/history/HistoryHandler.java +++ b/src-java/network-topology/network-storm-topology/src/main/java/org/openkilda/wfm/topology/network/storm/bolt/history/HistoryHandler.java @@ -24,20 +24,48 @@ import org.openkilda.wfm.topology.network.storm.bolt.history.command.HistoryCommand; import lombok.extern.slf4j.Slf4j; +import org.apache.storm.task.OutputCollector; +import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.tuple.Tuple; +import java.util.Map; + @Slf4j public class HistoryHandler extends AbstractBolt { private transient HistoryService historyService; + private static volatile PersistenceManager PERSISTENCE_MANAGER_INSTANCE; public HistoryHandler(PersistenceManager persistenceManager) { super(persistenceManager); } + /** + * Called when a task for this component is initialized within a worker on the cluster. + * It provides the bolt with the environment in which the bolt executes. + * Static instance of PersistenceManager is initialized at this step. + * + * @param stormConf The Storm configuration for this bolt. + * @param context This object can be used to get information about this task's place within the topology. + * @param collector The collector is used to emit tuples from this bolt. + */ + @Override + public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + synchronized (HistoryHandler.class) { + if (PERSISTENCE_MANAGER_INSTANCE == null) { + PERSISTENCE_MANAGER_INSTANCE = persistenceManager; + PERSISTENCE_MANAGER_INSTANCE.install(); + } + } + } + persistenceManager = null; + super.prepare(stormConf, context, collector); + } + @Override protected void init() { - this.historyService = new HistoryService(persistenceManager); + this.historyService = new HistoryService(PERSISTENCE_MANAGER_INSTANCE); } @Override