Class: TWPipeline::Stages::Segment

Inherits:
TWPipeline::Stage show all
Defined in:
lib/twpipeline/stages/segment.rb,
sig/twpipeline.rbs

Constant Summary collapse

CARRIED =

Returns:

  • (::Array[::Symbol])
%i[source host gate register].freeze

Constants inherited from TWPipeline::Stage

TWPipeline::Stage::REGISTRY

Instance Attribute Summary

Attributes inherited from TWPipeline::Stage

#options, #policy, #resources

Instance Method Summary collapse

Methods inherited from TWPipeline::Stage

#drop?, #flag, #initialize, lookup, order, ordered, register, #reject, slug

Constructor Details

This class inherits a constructor from TWPipeline::Stage

Instance Method Details

#call(input, output) ⇒ Object



8
9
10
11
12
13
14
15
16
17
# File 'lib/twpipeline/stages/segment.rb', line 8

def call(input, output)
  counts = Parallel.pipe(input, output, workers: resources.cores) do |line|
    stripped = line.strip
    next nil if stripped.empty?

    split(Jsonl.parse(stripped)).map { |record| Jsonl.dump(record) }
  end

  {documents: counts.fetch(:read), sentences: counts.fetch(:written)}
end