Module: Parallel
- Defined in:
- lib/parallel.rb,
lib/parallel/version.rb,
lib/parallel/serializer.rb
Defined Under Namespace
Modules: Serializer Classes: Break, DeadWorker, ExceptionWrapper, JobFactory, Kill, UndumpableException, UserInterruptHandler, Worker
Constant Summary collapse
- Stop =
Object.new.freeze
- VERSION =
rubocop:disable Naming/ConstantName
Version = '2.2.0'
Class Method Summary collapse
- .all?(*args, &block) ⇒ Boolean
- .any?(*args, &block) ⇒ Boolean
- .each(array, options = {}, &block) ⇒ Object
- .each_with_index(array, options = {}, &block) ⇒ Object
- .filter_map ⇒ Object
- .flat_map ⇒ Object
- .in_processes(options = {}, &block) ⇒ Object
- .in_threads(options = { count: 2 }) ⇒ Object
- .map(source, options = {}, &block) ⇒ Object
- .map_with_index(array, options = {}, &block) ⇒ Object
-
.physical_processor_count ⇒ Object
Number of physical processor cores on the current system.
-
.processor_count ⇒ Object
Number of processors seen by the OS or value considering CPU quota if the process is inside a cgroup, used for process scheduling.
- .worker_number ⇒ Object
-
.worker_number=(worker_num) ⇒ Object
TODO: this does not work when doing threads in forks, so should remove and yield the number instead if needed.
Class Method Details
.all?(*args, &block) ⇒ Boolean
256 257 258 259 |
# File 'lib/parallel.rb', line 256 def all?(*args, &block) raise "You must provide a block when calling #all?" if block.nil? !!each(*args) { |*a| raise Kill unless block.call(*a) } end |
.any?(*args, &block) ⇒ Boolean
251 252 253 254 |
# File 'lib/parallel.rb', line 251 def any?(*args, &block) raise "You must provide a block when calling #any?" if block.nil? !each(*args) { |*a| raise Kill if block.call(*a) } end |
.each(array, options = {}, &block) ⇒ Object
247 248 249 |
# File 'lib/parallel.rb', line 247 def each(array, = {}, &block) map(array, .merge(discard_results: true), &block) end |
.each_with_index(array, options = {}, &block) ⇒ Object
261 262 263 |
# File 'lib/parallel.rb', line 261 def each_with_index(array, = {}, &block) each(array, .merge(with_index: true), &block) end |
.filter_map ⇒ Object
333 334 335 |
# File 'lib/parallel.rb', line 333 def filter_map(...) map(...).select { |value| value } end |
.flat_map ⇒ Object
329 330 331 |
# File 'lib/parallel.rb', line 329 def flat_map(...) map(...).flatten(1) end |
.in_processes(options = {}, &block) ⇒ Object
241 242 243 244 245 |
# File 'lib/parallel.rb', line 241 def in_processes( = {}, &block) count, = () count ||= processor_count map(0...count, .merge(in_processes: count), &block) end |
.in_threads(options = { count: 2 }) ⇒ Object
225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 |
# File 'lib/parallel.rb', line 225 def in_threads( = { count: 2 }) threads = [] count, = () Thread.handle_interrupt(Exception => :never) do Thread.handle_interrupt(Exception => :immediate) do count.times do |i| threads << Thread.new { yield(i) } end threads.map(&:value) end ensure threads.each(&:kill) end end |
.map(source, options = {}, &block) ⇒ Object
265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 |
# File 'lib/parallel.rb', line 265 def map(source, = {}, &block) = .dup [:mutex] = Mutex.new if .slice(:in_processes, :in_threads, :in_ractors).size > 1 raise ArgumentError, "Use only one of `in_processes`, `in_threads`, or `in_ractors`." end if [:in_ractors] ? block : [:ractor] raise ArgumentError, "use either in_ractors with :ractor or a not in_ractors and block" end if RUBY_PLATFORM.include?('java') && ![:in_processes] method = :in_threads size = [method] || processor_count elsif [:in_threads] method = :in_threads size = [method] elsif [:in_ractors] method = :in_ractors size = [method] else method = :in_processes if Process.respond_to?(:fork) size = [method] || processor_count else warn "parallel: Process.fork is not supported by this Ruby" size = 0 end end raise ArgumentError, "worker count must be a non-negative Integer" unless size.is_a?(Integer) && size >= 0 job_factory = JobFactory.new(source, [:mutex]) size = [job_factory.size, size].min discard_results = [:discard_results] # finish callback needs the results, careful to do that before add_progress_bar which adds finish [:discard_results] = discard_results && ![:finish] (job_factory, ) result = if size == 0 block = ractor_block() if method == :in_ractors work_direct(job_factory, , &block) elsif method == :in_threads work_in_threads(job_factory, .merge(count: size), &block) elsif method == :in_ractors work_in_ractors(job_factory, .merge(count: size)) else work_in_processes(job_factory, .merge(count: size), &block) end return result.value if result.is_a?(Break) raise result if result.is_a?(Exception) discard_results ? source : result end |
.map_with_index(array, options = {}, &block) ⇒ Object
325 326 327 |
# File 'lib/parallel.rb', line 325 def map_with_index(array, = {}, &block) map(array, .merge(with_index: true), &block) end |
.physical_processor_count ⇒ Object
Number of physical processor cores on the current system.
338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 |
# File 'lib/parallel.rb', line 338 def physical_processor_count @physical_processor_count ||= begin ppc = case RbConfig::CONFIG["target_os"] when /darwin[12]/ IO.popen("/usr/sbin/sysctl -n hw.physicalcpu").read.to_i when /linux/ cores = {} # unique physical ID / core ID combinations phy = 0 File.read("/proc/cpuinfo").scan(/^physical id.*|^core id.*/) do |ln| if ln.start_with?("physical") phy = ln[/\d+/] elsif ln.start_with?("core") cid = "#{phy}:#{ln[/\d+/]}" cores[cid] = true unless cores[cid] end end cores.count when /mswin|mingw/ physical_processor_count_windows else processor_count end # fall back to logical count if physical info is invalid ppc > 0 ? ppc : processor_count end end |
.processor_count ⇒ Object
Number of processors seen by the OS or value considering CPU quota if the process is inside a cgroup, used for process scheduling
368 369 370 |
# File 'lib/parallel.rb', line 368 def processor_count @processor_count ||= Integer(ENV['PARALLEL_PROCESSOR_COUNT'] || available_processor_count) end |
.worker_number ⇒ Object
372 373 374 |
# File 'lib/parallel.rb', line 372 def worker_number Thread.current[:parallel_worker_number] end |
.worker_number=(worker_num) ⇒ Object
TODO: this does not work when doing threads in forks, so should remove and yield the number instead if needed
377 378 379 |
# File 'lib/parallel.rb', line 377 def worker_number=(worker_num) Thread.current[:parallel_worker_number] = worker_num end |