Class: Fluent::Plugin::OpentelemetryOutput::BatchProcessor
- Inherits:
-
Object
- Object
- Fluent::Plugin::OpentelemetryOutput::BatchProcessor
- Defined in:
- lib/fluent/plugin/out_opentelemetry.rb
Defined Under Namespace
Classes: InvalidRecords
Constant Summary collapse
- RESOURCE_KEY_MAP =
{ Opentelemetry::RECORD_TYPE_LOGS => "resourceLogs", Opentelemetry::RECORD_TYPE_METRICS => "resourceMetrics", Opentelemetry::RECORD_TYPE_TRACES => "resourceSpans" }.freeze
Class Method Summary collapse
Class Method Details
.build_export_requests(chunk, logger) ⇒ Object
134 135 136 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 |
# File 'lib/fluent/plugin/out_opentelemetry.rb', line 134 def self.build_export_requests(chunk, logger) requests = { Opentelemetry::RECORD_TYPE_LOGS => {}, Opentelemetry::RECORD_TYPE_METRICS => {}, Opentelemetry::RECORD_TYPE_TRACES => {} } invalid = Hash.new { |hash, key| hash[key] = InvalidRecords.new } chunk.each do |_, record| # rubocop:disable Style/HashEachMethods record_type = record["type"] resource_key = RESOURCE_KEY_MAP[record_type] unless resource_key invalid["unknown type"].add(record_type.inspect) next end begin record["message"] = JSON.parse(record["message"]) rescue JSON::ParserError, TypeError => e invalid["broken message"].add(e.) next end resources = record["message"][resource_key] if record["message"].is_a?(Hash) unless resources.is_a?(Array) && resources.first.is_a?(Hash) invalid["no #{resource_key}"].add next end resource_hash = resources.first["resource"].hash if requests[record_type][resource_hash].nil? requests[record_type][resource_hash] = record else requests[record_type][resource_hash]["message"][resource_key].concat(resources) end end unless invalid.empty? logger.warn do details = invalid.map { |reason, records| "#{reason}=#{records}" } "Skipped invalid records (total: #{invalid.values.sum(&:count)}): #{details.join(', ')}" end end merged_records = requests.values.flat_map(&:values) merged_records.each do |record| record["message"] = record["message"].to_json end merged_records end |