Class: Ingester
- Inherits:
-
Object
- Object
- Ingester
- Defined in:
- lib/fluent/plugin/ingester.rb
Class Attribute Summary collapse
-
.client_cache ⇒ Object
Returns the value of attribute client_cache.
Class Method Summary collapse
Instance Method Summary collapse
- #build_uri(container_sas_uri, name) ⇒ Object
-
#initialize(outconfiguration) ⇒ Ingester
constructor
A new instance of Ingester.
-
#post_message_to_queue_http(queue_uri_with_sas, message) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength.
-
#prepare_ingestion_message2(db, table, data_uri, blob_size_bytes, identity_token, compression_enabled = true, mapping_reference = nil) ⇒ Object
rubocop:disable Metrics/MethodLength.
-
#resources ⇒ Object
CRITICAL FIX: Dynamic resource access instead of stale cached reference.
- #token_provider ⇒ Object
-
#upload_data_to_blob_and_queue(raw_data, blob_name, db, table_name, compression_enabled = true, mapping_reference = nil) ⇒ Object
rubocop:enable Metrics/AbcSize, Metrics/MethodLength.
-
#upload_to_blob(blob_uri, raw_data, blob_name) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength.
Constructor Details
#initialize(outconfiguration) ⇒ Ingester
Returns a new instance of Ingester.
25 26 27 28 29 30 31 32 33 |
# File 'lib/fluent/plugin/ingester.rb', line 25 def initialize(outconfiguration) # Initialize Ingester with configuration and resources @client = self.class.client(outconfiguration) @logger = begin outconfiguration.logger rescue StandardError Logger.new($stdout) end end |
Class Attribute Details
.client_cache ⇒ Object
Returns the value of attribute client_cache.
22 23 24 |
# File 'lib/fluent/plugin/ingester.rb', line 22 def client_cache @client_cache end |
Class Method Details
.client(outconfiguration) ⇒ Object
35 36 37 38 39 40 41 42 43 44 |
# File 'lib/fluent/plugin/ingester.rb', line 35 def self.client(outconfiguration) # Thread-safe singleton client cache with basic validation return self.client_cache if self.client_cache # Double-checked locking pattern for thread safety @client_mutex ||= Mutex.new @client_mutex.synchronize do self.client_cache ||= Client.new(outconfiguration) end end |
Instance Method Details
#build_uri(container_sas_uri, name) ⇒ Object
51 52 53 54 55 56 |
# File 'lib/fluent/plugin/ingester.rb', line 51 def build_uri(container_sas_uri, name) # Build a blob URI with SAS token base_uri, sas_token = container_sas_uri.split('?', 2) base_uri = base_uri.chomp('/') "#{base_uri}/#{name}?#{sas_token}" end |
#post_message_to_queue_http(queue_uri_with_sas, message) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength
116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 |
# File 'lib/fluent/plugin/ingester.rb', line 116 def (queue_uri_with_sas, ) # Post the ingestion message to Azure Queue base_uri, sas_token = queue_uri_with_sas.split('?', 2) base_uri = base_uri.chomp('/') post_uri = URI("#{base_uri}/messages?#{sas_token}") = Base64.strict_encode64() request = Net::HTTP::Post.new(post_uri) request['Content-Type'] = 'application/xml' request.body = "<QueueMessage><MessageText>#{}</MessageText></QueueMessage>" response = Net::HTTP.start(post_uri.hostname, post_uri.port, use_ssl: post_uri.scheme == 'https', open_timeout: 10, read_timeout: 30, write_timeout: 10) do |http| http.request(request) end { code: response.code, message: response., body: response.body } end |
#prepare_ingestion_message2(db, table, data_uri, blob_size_bytes, identity_token, compression_enabled = true, mapping_reference = nil) ⇒ Object
rubocop:disable Metrics/MethodLength
91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 |
# File 'lib/fluent/plugin/ingester.rb', line 91 def (db, table, data_uri, blob_size_bytes, identity_token, compression_enabled = true, mapping_reference = nil) # Prepare the ingestion message for Azure Queue additional_props = { 'authorizationContext' => identity_token, 'format' => 'multijson' } additional_props['CompressionType'] = 'gzip' if compression_enabled additional_props['ingestionMappingReference'] = mapping_reference if mapping_reference && !mapping_reference.empty? { 'Id' => SecureRandom.uuid, 'BlobPath' => data_uri, 'RawDataSize' => blob_size_bytes, 'DatabaseName' => db, 'TableName' => table, 'RetainBlobOnSuccess' => true, 'FlushImmediately' => true, 'ReportLevel' => 2, # Report both failures and successes 'ReportMethod' => 0, # Use Azure Queue for reporting 'AdditionalProperties' => additional_props }.to_json end |
#resources ⇒ Object
CRITICAL FIX: Dynamic resource access instead of stale cached reference
47 48 49 |
# File 'lib/fluent/plugin/ingester.rb', line 47 def resources @client.resources end |
#token_provider ⇒ Object
148 149 150 151 |
# File 'lib/fluent/plugin/ingester.rb', line 148 def token_provider # Return the token provider from the client @client.token_provider end |
#upload_data_to_blob_and_queue(raw_data, blob_name, db, table_name, compression_enabled = true, mapping_reference = nil) ⇒ Object
rubocop:enable Metrics/AbcSize, Metrics/MethodLength
137 138 139 140 141 142 143 144 145 146 |
# File 'lib/fluent/plugin/ingester.rb', line 137 def upload_data_to_blob_and_queue(raw_data, blob_name, db, table_name, compression_enabled = true, mapping_reference = nil) # Upload data to blob and send ingestion message to queue # Use dynamic resources method instead of stale cached reference current_resources = resources blob_uri, blob_size_bytes = upload_to_blob(current_resources[:blob_sas_uri], raw_data, blob_name) = (db, table_name, blob_uri, blob_size_bytes, current_resources[:identity_token], compression_enabled, mapping_reference) (current_resources[:queue_sas_uri], ) { blob_uri: blob_uri, blob_size_bytes: blob_size_bytes } end |
#upload_to_blob(blob_uri, raw_data, blob_name) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength
59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 |
# File 'lib/fluent/plugin/ingester.rb', line 59 def upload_to_blob(blob_uri, raw_data, blob_name) # Upload raw data to Azure Blob Storage uri_str = build_uri(blob_uri, blob_name) uri = URI.parse(uri_str) blob_size = raw_data.bytesize request = Net::HTTP::Put.new(uri) request.body = raw_data request['x-ms-blob-type'] = 'BlockBlob' request['Content-Length'] = blob_size.to_s response = Net::HTTP.start(uri.hostname, uri.port, use_ssl: uri.scheme == 'https', open_timeout: 10, read_timeout: 30, write_timeout: 10) do |http| http.request(request) end unless response.code.to_i.between?(200, 299) begin error_handler = KustoErrorHandler.new(response.body) if error_handler.permanent_error? @logger.error("Permanent error while uploading blob: #{error_handler.}.Blob name: #{blob_name}") end rescue StandardError => e @logger.error("Failed to parse error response with KustoErrorHandler: #{e.}. Blob name: #{blob_name}") end raise "Blob upload failed: #{response.code} #{response.} - #{response.body}" end [uri.to_s, blob_size] end |