2121#include " DetectorsDCS/DataPointValue.h"
2222#include " DetectorsDCS/DeliveryType.h"
2323
24+ // #ifdef WITH_OPENMP
25+ // #include <omp.h>
26+ // #endif
27+
2428// / @brief Class to process DCS data points
2529
2630namespace o2
@@ -68,6 +72,9 @@ class DCSProcessor
6872 template <typename T>
6973 int processArrayType (const std::vector<DPID >& array, DeliveryType type, const std::unordered_map<DPID , DPVAL >& map, std::vector<uint64_t >& latestTimeStamp, std::unordered_map<DPID , T>& destmap);
7074
75+ template <typename T>
76+ void doSimpleMovingAverage (int nelements, std::deque<T>& vect, float & avg, bool & isSMA);
77+
7178 virtual void processChars ();
7279 virtual void processInts ();
7380 virtual void processDoubles ();
@@ -78,8 +85,6 @@ class DCSProcessor
7885 virtual void processBinaries ();
7986 virtual uint64_t processFlag (uint64_t flag, const char * alias);
8087
81- void doSimpleMovingAverage (int nelements, std::deque<int >& vect, float & avg, bool & isSMA);
82-
8388 DQChars& getVectorForAliasChar (const DPID & id) { return mDpscharsmap [id]; }
8489 DQInts& getVectorForAliasInt (const DPID & id) { return mDpsintsmap [id]; }
8590 DQDoubles& getVectorForAliasDouble (const DPID & id) { return mDpsdoublesmap [id]; }
@@ -88,9 +93,13 @@ class DCSProcessor
8893 DQStrings& getVectorForAliasString (const DPID & id) { return mDpsstringsmap [id]; }
8994 DQTimes& getVectorForAliasTime (const DPID & id) { return mDpstimesmap [id]; }
9095 DQBinaries& getVectorForAliasBinary (const DPID & id) { return mDpsbinariesmap [id]; }
91-
96+
97+ void setNThreads (int n);
98+ int getNThreads () const { return mNThreads ; }
99+
92100 private:
93- std::vector<float > mAvgTestInt ; // moving average for DP named TestInt0
101+ std::vector<float > mAvgTestInt ; // moving average for int DPs
102+ std::vector<float > mAvgTestDouble ; // moving average for double DPs
94103 std::unordered_map<DPID , DQChars> mDpscharsmap ;
95104 std::unordered_map<DPID , DQInts> mDpsintsmap ;
96105 std::unordered_map<DPID , DQDoubles> mDpsdoublesmap ;
@@ -115,7 +124,8 @@ class DCSProcessor
115124 std::vector<uint64_t > mLatestTimestampstrings ;
116125 std::vector<uint64_t > mLatestTimestamptimes ;
117126 std::vector<uint64_t > mLatestTimestampbinaries ;
118-
127+ int mNThreads = 1 ; // number of threads
128+
119129 ClassDefNV (DCSProcessor, 0 );
120130};
121131
@@ -140,6 +150,10 @@ int DCSProcessor::processArrayType(const std::vector<DPID>& array, DeliveryType
140150 int found = 0 ;
141151 auto s = array.size ();
142152 if (s > 0 ) {
153+ // #ifdef WITH_OPENMP
154+ // omp_set_num_threads(mNThreads);
155+ // #pragma omp parallel for schedule(dynamic)
156+ // #endif
143157 for (size_t i = 0 ; i != s; ++i) {
144158 auto it = processAlias (array[i], type, map);
145159 if (it == map.end ()) {
@@ -170,6 +184,26 @@ int DCSProcessor::processArrayType(const std::vector<DCSProcessor::DPID>& array,
170184template <>
171185int DCSProcessor::processArrayType (const std::vector<DCSProcessor::DPID >& array, DeliveryType type, const std::unordered_map<DCSProcessor::DPID , DCSProcessor::DPVAL >& map, std::vector<uint64_t >& latestTimeStamp, std::unordered_map<DCSProcessor::DPID , DCSProcessor::DQBinaries>& destmap);
172186
187+ template <typename T>
188+ void DCSProcessor::doSimpleMovingAverage (int nelements, std::deque<T>& vect, float & avg, bool & isSMA) {
189+
190+ // Do simple moving average on vector of type T
191+
192+ if (vect.size () < nelements) {
193+ avg += vect[vect.size () - 1 ];
194+ return ;
195+ }
196+ if (vect.size () == nelements) {
197+ avg += vect[vect.size () - 1 ];
198+ avg /= nelements;
199+ isSMA = true ;
200+ return ;
201+ }
202+ avg += (vect[vect.size () - 1 ] - vect[0 ]) / nelements;
203+ vect.pop_front ();
204+ isSMA = true ;
205+ }
206+
173207} // namespace dcs
174208} // namespace o2
175209
0 commit comments