Class: Fluent::BigQueryOutput

Inherits:
BufferedOutput
  • Object
show all
Defined in:
lib/fluent/plugin/out_bigquery.rb

Overview

TODO: error classes for each api error responses class BigQueryAPIError < StandardError end

Defined Under Namespace

Classes: BooleanFieldSchema, FieldSchema, FloatFieldSchema, IntegerFieldSchema, RecordSchema, StringFieldSchema, TimestampFieldSchema

Instance Method Summary collapse

Constructor Details

#initialize ⇒ BigQueryOutput

Table types https://developers.google.com/bigquery/docs/tables

type - The following data types are supported; see Data Formats for details on each data type: STRING INTEGER FLOAT BOOLEAN RECORD A JSON object, used when importing nested records. This type is only available when using JSON source files.

mode - Whether a field can be null. The following values are supported: NULLABLE - The cell can be null. REQUIRED - The cell cannot be null. REPEATED - Zero or more repeated simple or nested subfields. This mode is only supported when using JSON source files.



123
124
125
126
127
128
129
130
# File 'lib/fluent/plugin/out_bigquery.rb', line 123

def initialize
  super
  require 'json'
  require 'google/api_client'
  require 'google/api_client/client_secrets'
  require 'google/api_client/auth/installed_app'
  require 'google/api_client/auth/compute_service_account'
end

Instance Method Details

#client ⇒ Object



209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
# File 'lib/fluent/plugin/out_bigquery.rb', line 209

def client
  return @cached_client if @cached_client && @cached_client_expiration > Time.now

  client = Google::APIClient.new(
    application_name: 'Fluentd BigQuery plugin',
    application_version: Fluent::BigQueryPlugin::VERSION
  )

  case @auth_method
  when 'private_key'
    key = Google::APIClient::PKCS12.load_key( @private_key_path, @private_key_passphrase )
    asserter = Google::APIClient::JWTAsserter.new(
      @email,
      "https://www.googleapis.com/auth/bigquery",
      key
    )
    # refresh_auth
    client.authorization = asserter.authorize

  when 'compute_engine'
    auth = Google::APIClient::ComputeServiceAccount.new
    auth.fetch_access_token!
    client.authorization = auth

  else
    raise ConfigError, "Unknown auth method: #{@auth_method}"
  end

  @cached_client_expiration = Time.now + 1800
  @cached_client = client
end

#configure(conf) ⇒ Object



137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
# File 'lib/fluent/plugin/out_bigquery.rb', line 137

def configure(conf)
  super

  case @auth_method
  when 'private_key'
    unless @email && @private_key_path
      raise Fluent::ConfigError, "'email' and 'private_key_path' must be specified if auth_method == 'private_key'"
    end
  when 'compute_engine'
    # Do nothing
  else
    raise Fluent::ConfigError, "unrecognized 'auth_method': #{@auth_method}"
  end

  unless @table.nil? ^ @tables.nil?
    raise Fluent::ConfigError, "'table' or 'tables' must be specified, and both are invalid"
  end

  @tablelist = @tables ? @tables.split(',') : [@table]

  @fields = RecordSchema.new('record')
  if @schema_path
    @fields.load_schema(JSON.parse(File.read(@schema_path)))
  end

  types = %w(string integer float boolean timestamp)
  types.each do |type|
    raw_fields = instance_variable_get("@field_#{type}")
    next unless raw_fields
    raw_fields.split(',').each do |field|
      @fields.register_field field.strip, type.to_sym
    end
  end

  @localtime = false if @localtime.nil? && @utc

  @timef = TimeFormatter.new(@time_format, @localtime)

  if @time_field
    keys = @time_field.split('.')
    last_key = keys.pop
    @add_time_field = ->(record, time) {
      keys.inject(record) { |h, k| h[k] ||= {} }[last_key] = @timef.format(time)
      record
    }
  else
    @add_time_field = ->(record, time) { record }
  end

  if @insert_id_field
    insert_id_keys = @insert_id_field.split('.')
    @get_insert_id = ->(record) {
      insert_id_keys.inject(record) {|h, k| h[k] }
    }
  else
    @get_insert_id = nil
  end
end

#create_table(table_id) ⇒ Object



245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
# File 'lib/fluent/plugin/out_bigquery.rb', line 245

def create_table(table_id)
  res = client().execute(
    :api_method => @bq.tables.insert,
    :parameters => {
      'projectId' => @project,
      'datasetId' => @dataset,
    },
    :body_object => {
      'tableReference' => {
        'tableId' => table_id,
      },
      'schema' => {
        'fields' => @fields.to_a,
      },
    }
  )
  unless res.success?
    # api_error? -> client cache clear
    @cached_client = nil

    message = res.body
    if res.body =~ /^\{/
      begin
        res_obj = JSON.parse(res.body)
        message = res_obj['error']['message'] || res.body
      rescue => e
        log.warn "Parse error: google api error response body", :body => res.body
      end
      if res_obj and res_obj['code'] == 409 and /Already Exists:/ =~ message
        # ignore 'Already Exists' error
        return
      end
    end
    log.error "tables.insert API", :project_id => @project, :dataset => @dataset, :table => table_id, :code => res.status, :message => message
    raise "failed to create table in bigquery" # TODO: error class
  end
end

#extract_error_message(response_body) ⇒ Object



400
401
402
403
404
# File 'lib/fluent/plugin/out_bigquery.rb', line 400

def extract_error_message(response_body)
  res_obj = extract_response_obj(response_body)
  return response_body if res_obj.nil?
  res_obj['error']['message'] || response_body
end

#extract_response_obj(response_body) ⇒ Object

def client_oauth # not implemented raise NotImplementedError, "OAuth needs browser authentication..."

client = Google::APIClient.new(
application_name: 'Example Ruby application',
application_version: '1.0.0'
)
bigquery = client.discovered_api('bigquery', 'v2')
flow = Google::APIClient::InstalledAppFlow.new(
client_id: @client_id
client_secret: @client_secret
scope: ['https://www.googleapis.com/auth/bigquery']
)
client.authorization = flow.authorize # browser authentication !
client

end



392
393
394
395
396
397
398
# File 'lib/fluent/plugin/out_bigquery.rb', line 392

def extract_response_obj(response_body)
  return nil unless response_body =~ /^\{/
  JSON.parse(response_body)
rescue
  log.warn "Parse error: google api error response body", body: response_body
  return nil
end

#fetch_schema ⇒ Object



349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
# File 'lib/fluent/plugin/out_bigquery.rb', line 349

def fetch_schema
  table_id_format = @tablelist[0]
  table_id = generate_table_id(table_id_format, Time.at(Fluent::Engine.now))
  res = client.execute(
    api_method: @bq.tables.get,
    parameters: {
      'projectId' => @project,
      'datasetId' => @dataset,
      'tableId' => table_id,
    }
  )

  unless res.success?
    # api_error? -> client cache clear
    @cached_client = nil
    message = extract_error_message(res.body)
    log.error "tables.get API", project_id: @project, dataset: @dataset, table: table_id, code: res.status, message: message
    raise "failed to fetch schema from bigquery" # TODO: error class
  end

  res_obj = JSON.parse(res.body)
  schema = res_obj['schema']['fields']
  log.debug "Load schema from BigQuery: #{@project}:#{@dataset}.#{table_id} #{schema}"
  @fields.load_schema(schema, false)
end

#format_stream(tag, es) ⇒ Object



318
319
320
321
322
323
324
325
326
327
328
329
330
# File 'lib/fluent/plugin/out_bigquery.rb', line 318

def format_stream(tag, es)
  super
  buf = ''
  es.each do |time, record|
    row = @fields.format(@add_time_field.call(record, time))
    unless row.empty?
      row = {"json" => row}
      row['insertId'] = @get_insert_id.call(record) if @get_insert_id
      buf << row.to_msgpack
    end
  end
  buf
end

#generate_table_id(table_id_format, current_time) ⇒ Object



241
242
243
# File 'lib/fluent/plugin/out_bigquery.rb', line 241

def generate_table_id(table_id_format, current_time)
  current_time.strftime(table_id_format)
end

#insert(table_id_format, rows) ⇒ Object



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
# File 'lib/fluent/plugin/out_bigquery.rb', line 283

def insert(table_id_format, rows)
  table_id = generate_table_id(table_id_format, Time.at(Fluent::Engine.now))
  res = client().execute(
    api_method: @bq.tabledata.insert_all,
    parameters: {
      'projectId' => @project,
      'datasetId' => @dataset,
      'tableId' => table_id,
    },
    body_object: {
      "rows" => rows
    }
  )
  unless res.success?
    # api_error? -> client cache clear
    @cached_client = nil

    res_obj = extract_response_obj(res.body)
    message = res_obj['error']['message'] || res.body
    if res_obj
      if @auto_create_table and res_obj and res_obj['error']['code'] == 404 and /Not Found: Table/i =~ message.to_s
        # Table Not Found: Auto Create Table
        create_table(table_id)
      end
    end
    log.error "tabledata.insertAll API", project_id: @project, dataset: @dataset, table: table_id, code: res.status, message: message
    raise "failed to insert into bigquery" # TODO: error class
  end
end

#load ⇒ Object

Raises:

  • (NotImplementedError)


313
314
315
316
# File 'lib/fluent/plugin/out_bigquery.rb', line 313

def load
  # https://developers.google.com/bigquery/loading-data-into-bigquery#loaddatapostrequest
  raise NotImplementedError # TODO
end

#start ⇒ Object



196
197
198
199
200
201
202
203
204
205
206
207
# File 'lib/fluent/plugin/out_bigquery.rb', line 196

def start
  super

  @bq = client.discovered_api("bigquery", "v2") # TODO: refresh with specified expiration
  @cached_client = nil
  @cached_client_expiration = nil

  @tables_queue = @tablelist.dup.shuffle
  @tables_mutex = Mutex.new

  fetch_schema() if @fetch_schema
end

#write(chunk) ⇒ Object



332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
# File 'lib/fluent/plugin/out_bigquery.rb', line 332

def write(chunk)
  rows = []
  chunk.msgpack_each do |row_object|
    # TODO: row size limit
    rows << row_object
  end

  # TODO: method

  insert_table = @tables_mutex.synchronize do
    t = @tables_queue.shift
    @tables_queue.push t
    t
  end
  insert(insert_table, rows)
end