Eigen  5.0.1
 
Loading...
Searching...
No Matches
CoreThreadPoolDevice.h
1// This file is part of Eigen, a lightweight C++ template library
2// for linear algebra.
3//
4// Copyright (C) 2023 Charlie Schlosser <cs.schlosser@gmail.com>
5//
6// This Source Code Form is subject to the terms of the Mozilla
7// Public License v. 2.0. If a copy of the MPL was not distributed
8// with this file, You can obtain one at http://mozilla.org/MPL/2.0/.
9// SPDX-License-Identifier: MPL-2.0
10
11#ifndef EIGEN_CORE_THREAD_POOL_DEVICE_H
12#define EIGEN_CORE_THREAD_POOL_DEVICE_H
13
14namespace Eigen {
15
16// CoreThreadPoolDevice provides an easy-to-understand Device for parallelizing Eigen Core expressions with
17// Threadpool. Expressions are recursively split evenly until the evaluation cost is less than the threshold for
18// delegating the task to a thread.
19/*
20 a
21 / \
22 / \
23 / \
24 / \
25 / \
26 / \
27 / \
28 a e
29 / \ / \
30 / \ / \
31 / \ / \
32 a c e g
33 / \ / \ / \ / \
34 / \ / \ / \ / \
35 a b c d e f g h
36*/
37// Each task descends the binary tree to the left, delegates the right task to a new thread, and continues to the
38// left. This ensures that work is evenly distributed to the thread pool as quickly as possible and minimizes the number
39// of tasks created during the evaluation. Consider an expression that is divided into 8 chunks. The
40// primary task 'a' creates tasks 'e' 'c' and 'b', and executes its portion of the expression at the bottom of the
41// tree. Likewise, task 'e' creates tasks 'g' and 'f', and executes its portion of the expression.
42
43struct CoreThreadPoolDevice {
44 using Task = std::function<void()>;
45 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE CoreThreadPoolDevice(ThreadPool& pool, float threadCostThreshold = 3e-5f)
46 : m_pool(pool) {
47 eigen_assert(threadCostThreshold >= 0.0f && "threadCostThreshold must be non-negative");
48 m_costFactor = threadCostThreshold;
49 }
50
51 template <int PacketSize>
52 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE int calculateLevels(Index size, float cost) const {
53 eigen_assert(cost >= 0.0f && "cost must be non-negative");
54 Index numOps = size / PacketSize;
55 int actualThreads = numOps < m_pool.NumThreads() ? static_cast<int>(numOps) : m_pool.NumThreads();
56 float totalCost = static_cast<float>(numOps) * cost;
57 float idealThreads = totalCost * m_costFactor;
58 if (idealThreads < static_cast<float>(actualThreads)) {
59 idealThreads = numext::maxi(idealThreads, 1.0f);
60 actualThreads = numext::mini(actualThreads, static_cast<int>(idealThreads));
61 }
62 int maxLevel = internal::log2_ceil(actualThreads);
63 return maxLevel;
64 }
65
66// MSVC does not like inlining parallelForImpl
67#if EIGEN_COMP_MSVC && !EIGEN_COMP_CLANG
68#define EIGEN_PARALLEL_FOR_INLINE
69#else
70#define EIGEN_PARALLEL_FOR_INLINE EIGEN_STRONG_INLINE
71#endif
72
73 template <typename UnaryFunctor, int PacketSize>
74 EIGEN_DEVICE_FUNC EIGEN_PARALLEL_FOR_INLINE void parallelForImpl(Index begin, Index end, UnaryFunctor& f,
75 Barrier& barrier, int level) {
76 while (level > 0) {
77 level--;
78 Index size = end - begin;
79 eigen_assert(size % PacketSize == 0 && "this function assumes size is a multiple of PacketSize");
80 Index mid = begin + numext::round_down(size >> 1, PacketSize);
81 Task right = [this, mid, end, &f, &barrier, level]() {
82 parallelForImpl<UnaryFunctor, PacketSize>(mid, end, f, barrier, level);
83 };
84 m_pool.Schedule(std::move(right));
85 end = mid;
86 }
87 for (Index i = begin; i < end; i += PacketSize) f(i);
88 barrier.Notify();
89 }
90
91 template <typename BinaryFunctor, int PacketSize>
92 EIGEN_DEVICE_FUNC EIGEN_PARALLEL_FOR_INLINE void parallelForImpl(Index outerBegin, Index outerEnd, Index innerBegin,
93 Index innerEnd, BinaryFunctor& f, Barrier& barrier,
94 int level) {
95 while (level > 0) {
96 level--;
97 Index outerSize = outerEnd - outerBegin;
98 if (outerSize > 1) {
99 Index outerMid = outerBegin + (outerSize >> 1);
100 Task right = [this, &f, &barrier, outerMid, outerEnd, innerBegin, innerEnd, level]() {
101 parallelForImpl<BinaryFunctor, PacketSize>(outerMid, outerEnd, innerBegin, innerEnd, f, barrier, level);
102 };
103 m_pool.Schedule(std::move(right));
104 outerEnd = outerMid;
105 } else {
106 Index innerSize = innerEnd - innerBegin;
107 eigen_assert(innerSize % PacketSize == 0 && "this function assumes innerSize is a multiple of PacketSize");
108 Index innerMid = innerBegin + numext::round_down(innerSize >> 1, PacketSize);
109 Task right = [this, &f, &barrier, outerBegin, outerEnd, innerMid, innerEnd, level]() {
110 parallelForImpl<BinaryFunctor, PacketSize>(outerBegin, outerEnd, innerMid, innerEnd, f, barrier, level);
111 };
112 m_pool.Schedule(std::move(right));
113 innerEnd = innerMid;
114 }
115 }
116 for (Index outer = outerBegin; outer < outerEnd; outer++)
117 for (Index inner = innerBegin; inner < innerEnd; inner += PacketSize) f(outer, inner);
118 barrier.Notify();
119 }
120
121#undef EIGEN_PARALLEL_FOR_INLINE
122
123 template <typename UnaryFunctor, int PacketSize>
124 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void parallelFor(Index begin, Index end, UnaryFunctor& f, float cost) {
125 Index size = end - begin;
126 int maxLevel = calculateLevels<PacketSize>(size, cost);
127 Barrier barrier(1 << maxLevel);
128 parallelForImpl<UnaryFunctor, PacketSize>(begin, end, f, barrier, maxLevel);
129 barrier.Wait();
130 }
131
132 template <typename BinaryFunctor, int PacketSize>
133 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void parallelFor(Index outerBegin, Index outerEnd, Index innerBegin,
134 Index innerEnd, BinaryFunctor& f, float cost) {
135 Index outerSize = outerEnd - outerBegin;
136 Index innerSize = innerEnd - innerBegin;
137 Index size = outerSize * innerSize;
138 int maxLevel = calculateLevels<PacketSize>(size, cost);
139 Barrier barrier(1 << maxLevel);
140 parallelForImpl<BinaryFunctor, PacketSize>(outerBegin, outerEnd, innerBegin, innerEnd, f, barrier, maxLevel);
141 barrier.Wait();
142 }
143
144 ThreadPool& m_pool;
145 // costFactor is the cost of delegating a task to a thread
146 // the inverse is used to avoid a floating point division
147 float m_costFactor;
148};
149
150// specialization of coefficient-wise assignment loops for CoreThreadPoolDevice
151
152namespace internal {
153
154#ifdef EIGEN_PARSED_BY_DOXYGEN
155struct Kernel;
156#endif
157
158template <typename Kernel>
159struct cost_helper {
160 using SrcEvaluatorType = typename Kernel::SrcEvaluatorType;
161 using DstEvaluatorType = typename Kernel::DstEvaluatorType;
162 using SrcXprType = typename SrcEvaluatorType::XprType;
163 using DstXprType = typename DstEvaluatorType::XprType;
164 static constexpr Index Cost = functor_cost<SrcXprType>::Cost + functor_cost<DstXprType>::Cost;
165};
166
167template <typename Kernel>
168struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, DefaultTraversal, NoUnrolling> {
169 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost;
170 struct AssignmentFunctor : public Kernel {
171 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
172 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer, Index inner) {
173 this->assignCoeffByOuterInner(outer, inner);
174 }
175 };
176
177 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
178 const Index innerSize = kernel.innerSize();
179 const Index outerSize = kernel.outerSize();
180 constexpr float cost = static_cast<float>(XprEvaluationCost);
181 AssignmentFunctor functor(kernel);
182 device.template parallelFor<AssignmentFunctor, 1>(0, outerSize, 0, innerSize, functor, cost);
183 }
184};
185
186template <typename Kernel>
187struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, DefaultTraversal, InnerUnrolling> {
188 using DstXprType = typename Kernel::DstEvaluatorType::XprType;
189 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost, InnerSize = DstXprType::InnerSizeAtCompileTime;
190 struct AssignmentFunctor : public Kernel {
191 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
192 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer) {
193 copy_using_evaluator_DefaultTraversal_InnerUnrolling<Kernel, 0, InnerSize>::run(*this, outer);
194 }
195 };
196 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
197 const Index outerSize = kernel.outerSize();
198 AssignmentFunctor functor(kernel);
199 constexpr float cost = static_cast<float>(XprEvaluationCost) * static_cast<float>(InnerSize);
200 device.template parallelFor<AssignmentFunctor, 1>(0, outerSize, functor, cost);
201 }
202};
203
204template <typename Kernel>
205struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, InnerVectorizedTraversal, NoUnrolling> {
206 using PacketType = typename Kernel::PacketType;
207 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost, PacketSize = unpacket_traits<PacketType>::size,
208 SrcAlignment = Kernel::AssignmentTraits::SrcAlignment,
209 DstAlignment = Kernel::AssignmentTraits::DstAlignment;
210 struct AssignmentFunctor : public Kernel {
211 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
212 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer, Index inner) {
213 this->template assignPacketByOuterInner<Unaligned, Unaligned, PacketType>(outer, inner);
214 }
215 };
216 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
217 const Index innerSize = kernel.innerSize();
218 const Index outerSize = kernel.outerSize();
219 const float cost = static_cast<float>(XprEvaluationCost) * static_cast<float>(innerSize);
220 AssignmentFunctor functor(kernel);
221 device.template parallelFor<AssignmentFunctor, PacketSize>(0, outerSize, 0, innerSize, functor, cost);
222 }
223};
224
225template <typename Kernel>
226struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, InnerVectorizedTraversal, InnerUnrolling> {
227 using PacketType = typename Kernel::PacketType;
228 using DstXprType = typename Kernel::DstEvaluatorType::XprType;
229 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost, PacketSize = unpacket_traits<PacketType>::size,
230 SrcAlignment = Kernel::AssignmentTraits::SrcAlignment,
231 DstAlignment = Kernel::AssignmentTraits::DstAlignment,
232 InnerSize = DstXprType::InnerSizeAtCompileTime;
233 struct AssignmentFunctor : public Kernel {
234 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
235 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer) {
236 copy_using_evaluator_innervec_InnerUnrolling<Kernel, 0, InnerSize, SrcAlignment, DstAlignment>::run(*this, outer);
237 }
238 };
239 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
240 const Index outerSize = kernel.outerSize();
241 constexpr float cost = static_cast<float>(XprEvaluationCost) * static_cast<float>(InnerSize);
242 AssignmentFunctor functor(kernel);
243 device.template parallelFor<AssignmentFunctor, PacketSize>(0, outerSize, functor, cost);
244 }
245};
246
247template <typename Kernel>
248struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, SliceVectorizedTraversal, NoUnrolling> {
249 using Scalar = typename Kernel::Scalar;
250 using PacketType = typename Kernel::PacketType;
251 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost, PacketSize = unpacket_traits<PacketType>::size;
252 struct PacketAssignmentFunctor : public Kernel {
253 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE PacketAssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
254 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer, Index inner) {
255 this->template assignPacketByOuterInner<Unaligned, Unaligned, PacketType>(outer, inner);
256 }
257 };
258 struct ScalarAssignmentFunctor : public Kernel {
259 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE ScalarAssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
260 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index outer) {
261 const Index innerSize = this->innerSize();
262 const Index packetAccessSize = numext::round_down(innerSize, PacketSize);
263 for (Index inner = packetAccessSize; inner < innerSize; inner++) this->assignCoeffByOuterInner(outer, inner);
264 }
265 };
266 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
267 const Index outerSize = kernel.outerSize();
268 const Index innerSize = kernel.innerSize();
269 const Index packetAccessSize = numext::round_down(innerSize, PacketSize);
270 constexpr float packetCost = static_cast<float>(XprEvaluationCost);
271 const float scalarCost = static_cast<float>(XprEvaluationCost) * static_cast<float>(innerSize - packetAccessSize);
272 PacketAssignmentFunctor packetFunctor(kernel);
273 ScalarAssignmentFunctor scalarFunctor(kernel);
274 device.template parallelFor<PacketAssignmentFunctor, PacketSize>(0, outerSize, 0, packetAccessSize, packetFunctor,
275 packetCost);
276 device.template parallelFor<ScalarAssignmentFunctor, 1>(0, outerSize, scalarFunctor, scalarCost);
277 }
278};
279
280template <typename Kernel>
281struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, LinearTraversal, NoUnrolling> {
282 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost;
283 struct AssignmentFunctor : public Kernel {
284 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
285 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index index) { this->assignCoeff(index); }
286 };
287 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
288 const Index size = kernel.size();
289 constexpr float cost = static_cast<float>(XprEvaluationCost);
290 AssignmentFunctor functor(kernel);
291 device.template parallelFor<AssignmentFunctor, 1>(0, size, functor, cost);
292 }
293};
294
295template <typename Kernel>
296struct dense_assignment_loop_with_device<Kernel, CoreThreadPoolDevice, LinearVectorizedTraversal, NoUnrolling> {
297 using Scalar = typename Kernel::Scalar;
298 using PacketType = typename Kernel::PacketType;
299 static constexpr Index XprEvaluationCost = cost_helper<Kernel>::Cost,
300 RequestedAlignment = Kernel::AssignmentTraits::LinearRequiredAlignment,
301 PacketSize = unpacket_traits<PacketType>::size,
302 DstIsAligned = Kernel::AssignmentTraits::DstAlignment >= RequestedAlignment,
303 DstAlignment = packet_traits<Scalar>::AlignedOnScalar ? RequestedAlignment
304 : Kernel::AssignmentTraits::DstAlignment,
305 SrcAlignment = Kernel::AssignmentTraits::JointAlignment;
306 struct AssignmentFunctor : public Kernel {
307 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE AssignmentFunctor(Kernel& kernel) : Kernel(kernel) {}
308 EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void operator()(Index index) {
309 this->template assignPacket<DstAlignment, SrcAlignment, PacketType>(index);
310 }
311 };
312 static constexpr bool UsePacketSegment = Kernel::AssignmentTraits::UsePacketSegment;
313 using head_loop =
314 unaligned_dense_assignment_loop<PacketType, DstAlignment, SrcAlignment, UsePacketSegment, DstIsAligned>;
315 using tail_loop = unaligned_dense_assignment_loop<PacketType, DstAlignment, SrcAlignment, UsePacketSegment, false>;
316
317 static EIGEN_DEVICE_FUNC EIGEN_STRONG_INLINE void run(Kernel& kernel, CoreThreadPoolDevice& device) {
318 const Index size = kernel.size();
319 const Index alignedStart =
320 DstIsAligned ? 0 : internal::first_aligned<RequestedAlignment>(kernel.dstDataPtr(), size);
321 const Index alignedEnd = alignedStart + numext::round_down(size - alignedStart, PacketSize);
322
323 head_loop::run(kernel, 0, alignedStart);
324
325 constexpr float cost = static_cast<float>(XprEvaluationCost);
326 AssignmentFunctor functor(kernel);
327 device.template parallelFor<AssignmentFunctor, PacketSize>(alignedStart, alignedEnd, functor, cost);
328
329 tail_loop::run(kernel, alignedEnd, size);
330 }
331};
332
333} // namespace internal
334
335} // namespace Eigen
336
337#endif // EIGEN_CORE_THREAD_POOL_DEVICE_H