Module: TWPipeline::Sorting

Defined in:
lib/twpipeline/sorting.rb,
sig/twpipeline.rbs

Class Method Summary collapse

Instance Method Summary collapse

Class Method Details

.programObject



9
# File 'lib/twpipeline/sorting.rb', line 9

def program = @program ||= ENV["TWP_SORT"] || (system("command -v gsort > /dev/null 2>&1") ? "gsort" : "sort")

.sort(*extra, resources: TWPipeline.resources) ⇒ Object



11
12
13
14
# File 'lib/twpipeline/sorting.rb', line 11

def sort(*extra, resources: TWPipeline.resources)
  tmp = TWPipeline.work.join("tmp").tap(&:mkpath)
  [program, "-S", resources.sort_buffer, "--parallel", resources.cores.to_s, "-T", tmp.to_s, *extra].shelljoin
end

.tally(output_path, floor: 1, resources: TWPipeline.resources, collapse: false) ⇒ Object

Raises:



32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
# File 'lib/twpipeline/sorting.rb', line 32

def tally(output_path, floor: 1, resources: TWPipeline.resources, collapse: false)
  Pathname(output_path).dirname.mkpath
  command = "#{unique_pipeline(floor: floor, resources: resources, collapse: collapse)} > #{output_path.to_s.shellescape}"
  written = 0

  IO.popen(command, "w") do |sink|
    yield(
      -> (key) {
        sink.puts(key)
        written += 1
      }
    )
  end

  raise Error, "sort pipeline failed: #{command}" unless $CHILD_STATUS.success?

  {keys_written: written, rows: Pathname(output_path).each_line.count}
end

.unique_pipeline(floor:, resources:, collapse: false) ⇒ Object



16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
# File 'lib/twpipeline/sorting.rb', line 16

def unique_pipeline(floor:, resources:, collapse: false)
  head = if collapse
    ["LC_ALL=C #{sort("-u", resources: resources)}", "LC_ALL=C cut -f1"]
  else
    ["LC_ALL=C #{sort(resources: resources)}"]
  end

  [
    *head,
    "LC_ALL=C uniq -c",
    "LC_ALL=C sed -E 's/^ *([0-9]+) /\\1\t/'",
    "LC_ALL=C awk -F'\t' -v floor=#{floor.to_i} '$1 >= floor'",
    "LC_ALL=C #{sort("-t", "\t", "-k1,1nr", "-k2,2", resources: resources)}"
  ].join(" | ")
end

Instance Method Details

#self?.program::String

Returns:

  • (::String)


58
# File 'sig/twpipeline.rbs', line 58

def self?.program: () -> ::String

#self?.sort::String

Parameters:

Returns:

  • (::String)


59
# File 'sig/twpipeline.rbs', line 59

def self?.sort: (*::String extra, ?resources: Resources) -> ::String

#self?.tally {|arg0| ... } ⇒ ::Hash[::Symbol, untyped]

Parameters:

  • output_path (Object)
  • floor: (::Integer)
  • resources: (Resources)
  • collapse: (Boolean)

Yields:

Yield Parameters:

  • arg0 (^(::String) -> void)

Yield Returns:

  • (void)

Returns:

  • (::Hash[::Symbol, untyped])


61
# File 'sig/twpipeline.rbs', line 61

def self?.tally: (untyped output_path, ?floor: ::Integer, ?resources: Resources, ?collapse: bool) { (^(::String) -> void) -> void } -> ::Hash[::Symbol, untyped]

#self?.unique_pipeline::String

Parameters:

  • floor: (::Integer)
  • resources: (Resources)
  • collapse: (Boolean)

Returns:

  • (::String)


60
# File 'sig/twpipeline.rbs', line 60

def self?.unique_pipeline: (floor: ::Integer, resources: Resources, ?collapse: bool) -> ::String