Class: SplitPgDump::Worker

Inherits:
Object
  • Object
show all
Defined in:
lib/split_pgdump.rb

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize ⇒ Worker

Returns a new instance of Worker.



21
22
23
24
25
26
27
28
29
# File 'lib/split_pgdump.rb', line 21

def initialize
  @rules_file = 'split.rules'
  @output_file = 'dump.sql'
  @sorter = `which sort`.chomp
  @xargs = `which xargs`.chomp
  @rules = []
  @num_sorters = 0
  @could_fork = true
end

Instance Attribute Details

#could_fork ⇒ Object

Returns the value of attribute could_fork.



20
21
22
# File 'lib/split_pgdump.rb', line 20

def could_fork
  @could_fork
end

#num_sorters ⇒ Object

Returns the value of attribute num_sorters.



19
20
21
# File 'lib/split_pgdump.rb', line 19

def num_sorters
  @num_sorters
end

#output_file ⇒ Object

Returns the value of attribute output_file.



19
20
21
# File 'lib/split_pgdump.rb', line 19

def output_file
  @output_file
end

#rules ⇒ Object

Returns the value of attribute rules.



19
20
21
# File 'lib/split_pgdump.rb', line 19

def rules
  @rules
end

#rules_file ⇒ Object

Returns the value of attribute rules_file.



19
20
21
# File 'lib/split_pgdump.rb', line 19

def rules_file
  @rules_file
end

#sorter ⇒ Object

Returns the value of attribute sorter.



19
20
21
# File 'lib/split_pgdump.rb', line 19

def sorter
  @sorter
end

#xargs ⇒ Object

Returns the value of attribute xargs.



20
21
22
# File 'lib/split_pgdump.rb', line 20

def xargs
  @xargs
end

Instance Method Details

#clear_files ⇒ Object



35
36
37
38
39
# File 'lib/split_pgdump.rb', line 35

def clear_files
  FileUtils.rm_f output_file
  FileUtils.rm_rf Dir[File.join(tables_dir, '*')]
  FileUtils.mkdir_p tables_dir
end

#find_rule(table) ⇒ Object



55
56
57
# File 'lib/split_pgdump.rb', line 55

def find_rule(table)
  @rules.find{|rule| table =~ rule.regex}
end

#parse_rules ⇒ Object



41
42
43
44
45
46
47
48
49
50
51
52
53
# File 'lib/split_pgdump.rb', line 41

def parse_rules
  if File.exists?(rules_file)
    File.open(rules_file) do |f|
      f.each_line do |line|
        if rule = SplitPgDump::Rule.parse(line)
          @rules << rule
        end
      end
    end
  else
    puts "NO FILE #{rules_file}"  if $debug
  end
end

#process_copy_line(out, line) ⇒ Object



76
77
78
79
80
81
82
83
84
85
86
# File 'lib/split_pgdump.rb', line 76

def process_copy_line(out, line)
  if line =~ /^\\\.[\r\n]/
    @table.flush_all
    @table.copy_lines{|l| out.puts l}
    puts "Table #{@table.table} copied in \t#{"%.2f" % (Time.now - @start_time)}s" if $debug
    @table = nil
    @state = :schema
  else
    @table.add_line(line)
  end
end

#process_schema_line(out, line) ⇒ Object



59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
# File 'lib/split_pgdump.rb', line 59

def process_schema_line(out, line)
  if line =~ /^COPY (\w+) \(([^)]+)\) FROM stdin;/
    table_name, columns = $1, $2.split(', ')
    rule = find_rule("#@schema.#{table_name}")
    @table = SplitPgDump::Table.new(tables_dir, @schema, table_name, columns, rule)
    @tables << @table
    puts "Start to write table \t#{table_name}" if $debug
    @start_time = Time.now
    @state = :table
  else
    if line =~ /^SET search_path = ([^,]+)/
      @schema = $1
    end
    out.write line
  end
end

#sort_and_finish ⇒ Object



110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
# File 'lib/split_pgdump.rb', line 110

def sort_and_finish
  files = []
  for table in @tables
    for one_file in table.files.values
      sort_args = one_file.sort_args(table.sort_args).shelljoin
      files << [one_file, sort_args]
    end
  end
  unless @xargs.empty?
    num_sorters = [@num_sorters, 1].max
    xargs_cmd = [@xargs, '-L1', '-P', num_sorters.to_s, @sorter].shelljoin
    puts xargs_cmd  if $debug
    IO.popen(xargs_cmd, 'w+') do |io|
      files.each{|one_file, sort_args|
        puts sort_args  if $debug
        io.puts sort_args
      }
      io.close_write
      io.each_line{|l|
        puts l  if $debug
      }
    end
  else
    sorter = @sorter.shellescape
    commands = files.map{|one_file, sort_args| "#{sorter} #{sort_args}" }
    if @num_sorters > 1
      commands.each_slice(@num_sorters) do |cmd|
        cmd = cmd.map{|c| "{ #{c} & }"}  if @could_fork
        cmd = cmd.join(' ; ')
        cmd += ' ; wait '  if @could_fork
        puts cmd  if $debug
        system cmd
      end
    else
      commands.each do |cmd|
        puts cmd  if $debug
        system cmd
      end
    end
  end
  files.each{|one_file, sort_args| one_file.write_finish}
end

#tables_dir ⇒ Object



31
32
33
# File 'lib/split_pgdump.rb', line 31

def tables_dir
  output_file + '-tables'
end

#work(in_stream) ⇒ Object



88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
# File 'lib/split_pgdump.rb', line 88

def work(in_stream)
  @state = :schema
  @table = nil
  @tables = []
  @schema = 'public'

  File.open(output_file, 'w') do |out|
    in_stream.each_line do |line|
      case @state
      when :schema
        process_schema_line(out, line)
      when :table
        process_copy_line(out, line)
      end
    end
  end

  @start_time = Time.now
  sort_and_finish
  puts "Finished in #{Time.now - @start_time}s #{Process.pid}" if $debug
end