Module: Rimless::Extensions::Producer

Extended by:
ActiveSupport::Concern
Included in:
Rimless
Defined in:
lib/rimless/extensions/producer.rb

Overview

The top-level Apache Kafka message producer integration.

Class Method Summary collapse

Class Method Details

.async_message(data:, schema:, topic:, **args) ⇒ Object

Send a single message to Apache Kafka. The data is encoded according to the given Apache Avro schema. The destination Kafka topic may be a relative name, or a hash which is passed to the .topic method to manipulate the application details. The message is sent is an asynchronous, non-blocking way.

Parameters:

  • data (Hash{Symbol => Mixed}) —

    the raw data, unencoded

  • schema (String, Symbol) —

    the Apache Avro schema to use

  • topic (String, Symbol, Hash{Symbol => Mixed}) —

    the destination Apache Kafka topic

  • args (Hash{Symbol => Mixed}) —

    additional parameters, see: https://bit.ly/4tHjcVg



43
44
45
46
# File 'lib/rimless/extensions/producer.rb', line 43

def async_message(data:, schema:, topic:, **args)
  encoded = Rimless.encode(data, schema: schema)
  async_raw_message(data: encoded, topic: topic, **args)
end

.async_raw_message(data:, topic:, headers: nil, **args) ⇒ Object

Send a single message to Apache Kafka. The data is not touched, so you need to encode it yourself before you pass it in. The destination Kafka topic may be a relative name, or a hash which is passed to the .topic method to manipulate the application details. The message is sent is an asynchronous, non-blocking way.

Parameters:

  • data (Hash{Symbol => Mixed}) —

    the raw data, unencoded

  • topic (String, Symbol, Hash{Symbol => Mixed}) —

    the destination Apache Kafka topic

  • headers (Hash{String => String, Array<String>}, nil) (defaults to: nil) —

    the message headers to send

  • args (Hash{Symbol => Mixed}) —

    additional parameters, see: https://bit.ly/4tHjcVg



88
89
90
91
92
93
94
95
96
97
98
99
# File 'lib/rimless/extensions/producer.rb', line 88

def async_raw_message(data:, topic:, headers: nil, **args)
  args = args.merge(topic: topic(topic), payload: data)

  # A compatibility helper for headers, as WaterDrop is now more strict
  if headers.present?
    args[:headers] = headers
    args[:headers].deep_stringify_keys!.deep_transform_values!(&:to_s) \
      if headers.is_a? Hash
  end

  producer.produce_async(**args)
end

.sync_message(data:, schema:, topic:, **args) ⇒ Object Also known as: message

Send a single message to Apache Kafka. The data is encoded according to the given Apache Avro schema. The destination Kafka topic may be a relative name, or a hash which is passed to the .topic method to manipulate the application details. The message is sent is a synchronous, blocking way.

Parameters:

  • data (Hash{Symbol => Mixed}) —

    the raw data, unencoded

  • schema (String, Symbol) —

    the Apache Avro schema to use

  • topic (String, Symbol, Hash{Symbol => Mixed}) —

    the destination Apache Kafka topic

  • args (Hash{Symbol => Mixed}) —

    additional parameters, see: https://bit.ly/4tHjcVg



25
26
27
28
# File 'lib/rimless/extensions/producer.rb', line 25

def sync_message(data:, schema:, topic:, **args)
  encoded = Rimless.encode(data, schema: schema)
  sync_raw_message(data: encoded, topic: topic, **args)
end

.sync_raw_message(data:, topic:, headers: nil, **args) ⇒ Object Also known as: raw_message

Send a single message to Apache Kafka. The data is not transformed, so you need to encode it yourself before you pass it in. The destination Kafka topic may be a relative name, or a hash which is passed to the .topic method to manipulate the application details. The message is sent is a synchronous, blocking way.

Parameters:

  • data (Hash{Symbol => Mixed}) —

    the raw data, unencoded

  • topic (String, Symbol, Hash{Symbol => Mixed}) —

    the destination Apache Kafka topic

  • headers (Hash{String => String, Array<String>}, nil) (defaults to: nil) —

    the message headers to send

  • args (Hash{Symbol => Mixed}) —

    additional parameters, see: https://bit.ly/4tHjcVg



61
62
63
64
65
66
67
68
69
70
71
72
# File 'lib/rimless/extensions/producer.rb', line 61

def sync_raw_message(data:, topic:, headers: nil, **args)
  args = args.merge(topic: topic(topic), payload: data)

  # A compatibility helper for headers, as WaterDrop is now more strict
  if headers.present?
    args[:headers] = headers
    args[:headers].deep_stringify_keys!.deep_transform_values!(&:to_s) \
      if headers.is_a? Hash
  end

  producer.produce_sync(**args)
end