Class: Fluent::GoogleCloudOutput
- Inherits:
-
BufferedOutput
- Object
- BufferedOutput
- Fluent::GoogleCloudOutput
- Defined in:
- lib/fluent/plugin/out_google_cloud.rb
Constant Summary collapse
- APPENGINE_SERVICE =
Constants for Google service names.
'appengine.googleapis.com'- COMPUTE_SERVICE =
'compute.googleapis.com'- DATAFLOW_SERVICE =
'dataflow.googleapis.com'
Instance Attribute Summary collapse
-
#common_labels ⇒ Object
readonly
Returns the value of attribute common_labels.
-
#gae_backend_name ⇒ Object
readonly
Returns the value of attribute gae_backend_name.
-
#gae_backend_version ⇒ Object
readonly
Returns the value of attribute gae_backend_version.
-
#project_id ⇒ Object
readonly
Expose attr_readers to make testing of metadata more direct than only testing it indirectly through metadata sent with logs.
-
#running_on_managed_vm ⇒ Object
readonly
Returns the value of attribute running_on_managed_vm.
-
#service_name ⇒ Object
readonly
Returns the value of attribute service_name.
-
#vm_id ⇒ Object
readonly
Returns the value of attribute vm_id.
-
#zone ⇒ Object
readonly
Returns the value of attribute zone.
Instance Method Summary collapse
- #configure(conf) ⇒ Object
- #format(tag, time, record) ⇒ Object
-
#initialize ⇒ GoogleCloudOutput
constructor
A new instance of GoogleCloudOutput.
- #shutdown ⇒ Object
- #start ⇒ Object
- #write(chunk) ⇒ Object
Constructor Details
#initialize ⇒ GoogleCloudOutput
Returns a new instance of GoogleCloudOutput.
52 53 54 55 56 57 58 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 52 def initialize super require 'cgi' require 'google/api_client' require 'google/api_client/auth/compute_service_account' require 'open-uri' end |
Instance Attribute Details
#common_labels ⇒ Object (readonly)
Returns the value of attribute common_labels.
50 51 52 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 50 def common_labels @common_labels end |
#gae_backend_name ⇒ Object (readonly)
Returns the value of attribute gae_backend_name.
47 48 49 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 47 def gae_backend_name @gae_backend_name end |
#gae_backend_version ⇒ Object (readonly)
Returns the value of attribute gae_backend_version.
48 49 50 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 48 def gae_backend_version @gae_backend_version end |
#project_id ⇒ Object (readonly)
Expose attr_readers to make testing of metadata more direct than only testing it indirectly through metadata sent with logs.
43 44 45 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 43 def project_id @project_id end |
#running_on_managed_vm ⇒ Object (readonly)
Returns the value of attribute running_on_managed_vm.
46 47 48 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 46 def running_on_managed_vm @running_on_managed_vm end |
#service_name ⇒ Object (readonly)
Returns the value of attribute service_name.
49 50 51 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 49 def service_name @service_name end |
#vm_id ⇒ Object (readonly)
Returns the value of attribute vm_id.
45 46 47 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 45 def vm_id @vm_id end |
#zone ⇒ Object (readonly)
Returns the value of attribute zone.
44 45 46 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 44 def zone @zone end |
Instance Method Details
#configure(conf) ⇒ Object
60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 60 def configure(conf) super case @auth_method when 'private_key' if !@private_key_email raise Fluent::ConfigError, ('"private_key_email" must be ' + 'specified if auth_method is "private_key"') elsif !@private_key_path raise Fluent::ConfigError, ('"private_key_path" must be ' + 'specified if auth_method is "private_key"') elsif !@private_key_passphrase raise Fluent::ConfigError, ('"private_key_passphrase" must be ' + 'specified if auth_method is "private_key"') end when 'compute_engine_service_account' # pass else raise Fluent::ConfigError, ('Unrecognized "auth_method" parameter. Please specify either ' + '"compute_engine_service_account" or "private_key".') end end |
#format(tag, time, record) ⇒ Object
134 135 136 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 134 def format(tag, time, record) [tag, time, record].to_msgpack end |
#shutdown ⇒ Object
130 131 132 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 130 def shutdown super end |
#start ⇒ Object
84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 84 def start super init_api_client() @successful_call = false # Grab metadata about the Google Compute Engine instance that we're on. @project_id = ('project/project-id') fully_qualified_zone = ('instance/zone') @zone = fully_qualified_zone.rpartition('/')[2] @vm_id = ('instance/id') # TODO: Send instance tags and/or hostname with the logs as well? @common_labels = {} # If this is running on a Managed VM, grab the relevant App Engine # metadata as well. # TODO: Add config options for these to allow for running outside GCE? attributes_string = ('instance/attributes/') attributes = attributes_string.split if (attributes.include?('gae_backend_name') && attributes.include?('gae_backend_version')) @running_on_managed_vm = true @gae_backend_name = ('instance/attributes/gae_backend_name') @gae_backend_version = ('instance/attributes/gae_backend_version') @service_name = APPENGINE_SERVICE common_labels["#{APPENGINE_SERVICE}/module_id"] = @gae_backend_name common_labels["#{APPENGINE_SERVICE}/version_id"] = @gae_backend_version elsif (attributes.include?('job_id')) @running_on_managed_vm = false @service_name = DATAFLOW_SERVICE @dataflow_job_id = ('instance/attributes/job_id') common_labels["#{DATAFLOW_SERVICE}/job_id"] = @dataflow_job_id else @running_on_managed_vm = false @service_name = COMPUTE_SERVICE end if (@service_name != DATAFLOW_SERVICE) common_labels["#{COMPUTE_SERVICE}/resource_type"] = 'instance' common_labels["#{COMPUTE_SERVICE}/resource_id"] = @vm_id end end |
#write(chunk) ⇒ Object
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 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 |
# File 'lib/fluent/plugin/out_google_cloud.rb', line 138 def write(chunk) # Group the entries since we have to make one call per tag. grouped_entries = {} chunk.msgpack_each do |tag, *arr| if !grouped_entries.has_key?(tag) grouped_entries[tag] = [] end grouped_entries[tag].push(arr) end grouped_entries.each do |tag, arr| write_log_entries_request = { 'commonLabels' => @common_labels, 'entries' => [], } arr.each do |time, record| if (record.has_key?('timeNanos')) ts_secs = (record['timeNanos'] / 1000000000).to_i ts_nanos = record['timeNanos'] % 1000000000 record.delete('timeNanos') else = Time.at(time) ts_secs = .tv_sec ts_nanos = .tv_nsec end entry = { 'metadata' => { 'serviceName' => @service_name, 'projectId' => @project_id, 'zone' => @zone, 'timestamp' => { 'seconds' => ts_secs, 'nanos' => ts_nanos }, }, } if record.has_key?('severity') entry['metadata']['severity'] = parse_severity(record['severity']) record.delete('severity') else entry['metadata']['severity'] = 'DEFAULT' end # use textPayload if the only remainaing key is 'message', # otherwise use a struct. if (record.size == 1 && record.has_key?('message')) entry['textPayload'] = record['message'] else entry['structPayload'] = record end write_log_entries_request['entries'].push(entry) end # Add a prefix to VMEngines logs to prevent namespace collisions, # and also escape the log name. log_name = CGI::escape(@running_on_managed_vm ? "#{APPENGINE_SERVICE}/#{tag}" : tag) url = ('https://logging.googleapis.com/v1beta3/projects/' + "#{@project_id}/logs/#{log_name}/entries:write") begin client = api_client() request = client.generate_request({ :uri => url, :body_object => write_log_entries_request, :http_method => 'POST', :authenticated => true }) client.execute!(request) # Let the user explicitly know when the first call succeeded, # to aid with verification and troubleshooting. if (!@successful_call) @successful_call = true $log.info "Successfully sent to Google Cloud Logging API." end # Allow most exceptions to propagate, which will cause fluentd to # retry (with backoff), but in some cases we catch the error and # drop the request (we will emit a log message in those cases). rescue Google::APIClient::ClientError => error # Most ClientErrors indicate a problem with the request itself and # should not be retried, unless it is an authentication issue, in # which case we will retry the request via re-raising the exception. if (is_retriable_client_error(error)) raise error end log_write_failure(write_log_entries_request, error) rescue JSON::GeneratorError => error # This happens if the request contains illegal characters; # do not retry it because it will fail repeatedly. log_write_failure(write_log_entries_request, error) end end end |