Class: Embulk::Output::Vertica::OutputThread

Inherits:
Object
  • Object
show all
Defined in:
lib/embulk/output/vertica/output_thread.rb

Instance Method Summary collapse

Constructor Details

#initialize(task) ⇒ OutputThread

Returns a new instance of OutputThread.



48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
# File 'lib/embulk/output/vertica/output_thread.rb', line 48

def initialize(task)
  @task = task
  @queue = SizedQueue.new(1)
  @num_input_rows = 0
  @num_output_rows = 0
  @num_rejected_rows = 0
  @outer_thread = Thread.current
  @thread_active = false
  @progress_log_timer = Time.now
  @previous_num_input_rows = 0

  case task['compress']
  when 'GZIP'
    @write_proc = self.method(:write_gzip)
  else
    @write_proc = self.method(:write_uncompressed)
  end
end

Instance Method Details

#abort_on_errorObject



211
212
213
# File 'lib/embulk/output/vertica/output_thread.rb', line 211

def abort_on_error
  @task['abort_on_error'] ? ' ABORT ON ERROR' : ''
end

#commitObject



162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
# File 'lib/embulk/output/vertica/output_thread.rb', line 162

def commit
  @thread_active = false
  if @thread.alive?
    Embulk.logger.debug { "embulk-output-vertica: push finish" }
    @queue.push('finish')
    Thread.pass
    @thread.join
  else
    raise RuntimeError, "embulk-output-vertica: thread died accidently"
  end

  task_report = {
    'num_input_rows' => @num_input_rows,
    'num_output_rows' => @num_output_rows,
    'num_rejected_rows' => @num_rejected_rows,
  }
end

#compressObject



203
204
205
# File 'lib/embulk/output/vertica/output_thread.rb', line 203

def compress
  " #{@task['compress']}"
end

#copy(conn, sql, &block) ⇒ Object

private



182
183
184
185
# File 'lib/embulk/output/vertica/output_thread.rb', line 182

def copy(conn, sql, &block)
  Embulk.logger.debug "embulk-output-vertica: #{sql}"
  results, rejects = conn.copy(sql, &block)
end

#copy_modeObject



207
208
209
# File 'lib/embulk/output/vertica/output_thread.rb', line 207

def copy_mode
  " #{@task['copy_mode']}"
end

#copy_sqlObject



187
188
189
# File 'lib/embulk/output/vertica/output_thread.rb', line 187

def copy_sql
  @copy_sql ||= "COPY #{quoted_schema}.#{quoted_temp_table} FROM STDIN#{compress}#{fjsonparser}#{copy_mode}#{abort_on_error} NO COMMIT"
end

#enqueue(json_page) ⇒ Object



67
68
69
70
71
72
73
74
75
# File 'lib/embulk/output/vertica/output_thread.rb', line 67

def enqueue(json_page)
  if @thread_active and @thread.alive?
    Embulk.logger.trace { "embulk-output-vertica: enqueue" }
    @queue.push(json_page)
  else
    Embulk.logger.info { "embulk-output-vertica: thread is dead, but still trying to enqueue" }
    raise RuntimeError, "embulk-output-vertica: thread is died, but still trying to enqueue"
  end
end

#fjsonparserObject



215
216
217
# File 'lib/embulk/output/vertica/output_thread.rb', line 215

def fjsonparser
  " PARSER fjsonparser(#{reject_on_materialized_type_error})"
end

#num_format(number) ⇒ Object



105
106
107
# File 'lib/embulk/output/vertica/output_thread.rb', line 105

def num_format(number)
  number.to_s.gsub(/(\d)(?=(\d{3})+(?!\d))/, '\1,')
end

#quoted_schemaObject



191
192
193
# File 'lib/embulk/output/vertica/output_thread.rb', line 191

def quoted_schema
  ::Jvertica.quote_identifier(@task['schema'])
end

#quoted_tableObject



195
196
197
# File 'lib/embulk/output/vertica/output_thread.rb', line 195

def quoted_table
  ::Jvertica.quote_identifier(@task['table'])
end

#quoted_temp_tableObject



199
200
201
# File 'lib/embulk/output/vertica/output_thread.rb', line 199

def quoted_temp_table
  ::Jvertica.quote_identifier(@task['temp_table'])
end

#reject_on_materialized_type_errorObject



219
220
221
# File 'lib/embulk/output/vertica/output_thread.rb', line 219

def reject_on_materialized_type_error
  @task['reject_on_materialized_type_error'] ? 'reject_on_materialized_type_error=true' : ''
end

#runObject



109
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
152
153
154
155
# File 'lib/embulk/output/vertica/output_thread.rb', line 109

def run
  Embulk.logger.debug { "embulk-output-vertica: thread started" }
  Vertica.connect(@task) do |jv|
    begin
      last_record = nil
      num_output_rows, rejects = copy(jv, copy_sql) do |stdin|
        while json_page = @queue.pop
          if json_page == 'finish'
            Embulk.logger.debug { "embulk-output-vertica: popped finish" }
            break
          end
          Embulk.logger.trace { "embulk-output-vertica: dequeued" }

          @write_proc.call(stdin, json_page) do |record|
            last_record = record
          end
        end
      end
      Embulk.logger.debug { "embulk-output-vertica: thread finished" }
      num_rejected_rows = rejects.size
      @num_output_rows += num_output_rows
      @num_rejected_rows += num_rejected_rows
      Embulk.logger.info { "embulk-output-vertica: COMMIT!" }
      jv.commit
      Embulk.logger.debug { "embulk-output-vertica: COMMITTED!" }
    rescue java.sql.SQLDataException => e
      if @task['reject_on_materialized_type_error'] and e.message =~ /Rejected by user-defined parser/
        Embulk.logger.warn "embulk-output-vertica: ROLLBACK! some of column types and values types do not fit #{last_record}"
      else
        Embulk.logger.warn "embulk-output-vertica: ROLLBACK!"
      end
      Embulk.logger.info { "embulk-output-vertica: last_record: #{last_record}" }
      jv.rollback
      raise e # die transaction
    rescue => e
      Embulk.logger.warn "embulk-output-vertica: ROLLBACK!"
      jv.rollback
      raise e
    end
  end
rescue => e
  @thread_active = false # not to be enqueued any more
  while @queue.size > 0
    @queue.pop # dequeue all because some might be still trying @queue.push and get blocked, need to release
  end
  @outer_thread.raise e.class.new("#{e.message}\n  #{e.backtrace.join("\n  ")}")
end

#startObject



157
158
159
160
# File 'lib/embulk/output/vertica/output_thread.rb', line 157

def start
  @thread = Thread.new(&method(:run))
  @thread_active = true
end

#write_buf(buf, json_page, &block) ⇒ Object



89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
# File 'lib/embulk/output/vertica/output_thread.rb', line 89

def write_buf(buf, json_page, &block)
  json_page.each do |record|
    yield(record) if block_given?
    Embulk.logger.trace { "embulk-output-vertica: record #{record}" }
    buf << record << "\n"
    @num_input_rows += 1
  end
  now = Time.now
  if @progress_log_timer < now - 10 # once in 10 seconds
    speed = ((@num_input_rows - @previous_num_input_rows) / (now - @progress_log_timer).to_f).round(1)
    @progress_log_timer = now
    @previous_num_input_rows = @num_input_rows
    Embulk.logger.info { "embulk-output-vertica: num_input_rows #{num_format(@num_input_rows)} (#{num_format(speed)} rows/sec)" }
  end
end

#write_gzip(io, page, &block) ⇒ Object



77
78
79
80
81
# File 'lib/embulk/output/vertica/output_thread.rb', line 77

def write_gzip(io, page, &block)
  buf = Zlib::Deflate.new
  write_buf(buf, page, &block)
  io << buf.finish
end

#write_uncompressed(io, page, &block) ⇒ Object



83
84
85
86
87
# File 'lib/embulk/output/vertica/output_thread.rb', line 83

def write_uncompressed(io, page, &block)
  buf = ''
  write_buf(buf, page, &block)
  io << buf
end