Class: Kafka::Broker

Inherits:
Object
  • Object
show all
Defined in:
lib/kafka/broker.rb

Instance Method Summary collapse

Constructor Details

#initialize(connection_builder:, host:, port:, node_id: nil, logger:) ⇒ Broker

Returns a new instance of Broker.



9
10
11
12
13
14
15
16
# File 'lib/kafka/broker.rb', line 9

def initialize(connection_builder:, host:, port:, node_id: nil, logger:)
  @connection_builder = connection_builder
  @connection = nil
  @host = host
  @port = port
  @node_id = node_id
  @logger = TaggedLogger.new(logger)
end

Instance Method Details

#add_offsets_to_txn(**options) ⇒ Object



185
186
187
188
189
# File 'lib/kafka/broker.rb', line 185

def add_offsets_to_txn(**options)
  request = Protocol::AddOffsetsToTxnRequest.new(**options)

  send_request(request)
end

#add_partitions_to_txn(**options) ⇒ Object



173
174
175
176
177
# File 'lib/kafka/broker.rb', line 173

def add_partitions_to_txn(**options)
  request = Protocol::AddPartitionsToTxnRequest.new(**options)

  send_request(request)
end

#address_match?(host, port) ⇒ Boolean

Returns:

  • (Boolean)


18
19
20
# File 'lib/kafka/broker.rb', line 18

def address_match?(host, port)
  host == @host && port == @port
end

#alter_configs(**options) ⇒ Object



137
138
139
140
141
# File 'lib/kafka/broker.rb', line 137

def alter_configs(**options)
  request = Protocol::AlterConfigsRequest.new(**options)

  send_request(request)
end

#api_versionsObject



155
156
157
158
159
# File 'lib/kafka/broker.rb', line 155

def api_versions
  request = Protocol::ApiVersionsRequest.new

  send_request(request)
end

#commit_offsets(**options) ⇒ Object



83
84
85
86
87
# File 'lib/kafka/broker.rb', line 83

def commit_offsets(**options)
  request = Protocol::OffsetCommitRequest.new(**options)

  send_request(request)
end

#connected?Boolean

Returns:

  • (Boolean)


33
34
35
# File 'lib/kafka/broker.rb', line 33

def connected?
  !@connection.nil?
end

#create_partitions(**options) ⇒ Object



143
144
145
146
147
# File 'lib/kafka/broker.rb', line 143

def create_partitions(**options)
  request = Protocol::CreatePartitionsRequest.new(**options)

  send_request(request)
end

#create_topics(**options) ⇒ Object



119
120
121
122
123
# File 'lib/kafka/broker.rb', line 119

def create_topics(**options)
  request = Protocol::CreateTopicsRequest.new(**options)

  send_request(request)
end

#delete_topics(**options) ⇒ Object



125
126
127
128
129
# File 'lib/kafka/broker.rb', line 125

def delete_topics(**options)
  request = Protocol::DeleteTopicsRequest.new(**options)

  send_request(request)
end

#describe_configs(**options) ⇒ Object



131
132
133
134
135
# File 'lib/kafka/broker.rb', line 131

def describe_configs(**options)
  request = Protocol::DescribeConfigsRequest.new(**options)

  send_request(request)
end

#describe_groups(**options) ⇒ Object



161
162
163
164
165
# File 'lib/kafka/broker.rb', line 161

def describe_groups(**options)
  request = Protocol::DescribeGroupsRequest.new(**options)

  send_request(request)
end

#disconnectnil

Returns:

  • (nil)


28
29
30
# File 'lib/kafka/broker.rb', line 28

def disconnect
  connection.close if connected?
end

#end_txn(**options) ⇒ Object



179
180
181
182
183
# File 'lib/kafka/broker.rb', line 179

def end_txn(**options)
  request = Protocol::EndTxnRequest.new(**options)

  send_request(request)
end

#fetch_messages(**options) ⇒ Kafka::Protocol::FetchResponse

Fetches messages from a specified topic and partition.

Parameters:

  • max_wait_time (Integer)
  • min_bytes (Integer)
  • topics (Hash)

Returns:



51
52
53
54
55
# File 'lib/kafka/broker.rb', line 51

def fetch_messages(**options)
  request = Protocol::FetchRequest.new(**options)

  send_request(request)
end

#fetch_metadata(**options) ⇒ Kafka::Protocol::MetadataResponse

Fetches cluster metadata from the broker.

Parameters:

  • topics (Array<String>)

Returns:



41
42
43
44
45
# File 'lib/kafka/broker.rb', line 41

def (**options)
  request = Protocol::MetadataRequest.new(**options)

  send_request(request)
end

#fetch_offsets(**options) ⇒ Object



77
78
79
80
81
# File 'lib/kafka/broker.rb', line 77

def fetch_offsets(**options)
  request = Protocol::OffsetFetchRequest.new(**options)

  send_request(request)
end

#find_coordinator(**options) ⇒ Object



107
108
109
110
111
# File 'lib/kafka/broker.rb', line 107

def find_coordinator(**options)
  request = Protocol::FindCoordinatorRequest.new(**options)

  send_request(request)
end

#heartbeat(**options) ⇒ Object



113
114
115
116
117
# File 'lib/kafka/broker.rb', line 113

def heartbeat(**options)
  request = Protocol::HeartbeatRequest.new(**options)

  send_request(request)
end

#init_producer_id(**options) ⇒ Object



167
168
169
170
171
# File 'lib/kafka/broker.rb', line 167

def init_producer_id(**options)
  request = Protocol::InitProducerIDRequest.new(**options)

  send_request(request)
end

#join_group(**options) ⇒ Object



89
90
91
92
93
# File 'lib/kafka/broker.rb', line 89

def join_group(**options)
  request = Protocol::JoinGroupRequest.new(**options)

  send_request(request)
end

#leave_group(**options) ⇒ Object



101
102
103
104
105
# File 'lib/kafka/broker.rb', line 101

def leave_group(**options)
  request = Protocol::LeaveGroupRequest.new(**options)

  send_request(request)
end

#list_groupsObject



149
150
151
152
153
# File 'lib/kafka/broker.rb', line 149

def list_groups
  request = Protocol::ListGroupsRequest.new

  send_request(request)
end

#list_offsets(**options) ⇒ Kafka::Protocol::ListOffsetResponse

Lists the offset of the specified topics and partitions.

Parameters:

  • topics (Hash)

Returns:



61
62
63
64
65
# File 'lib/kafka/broker.rb', line 61

def list_offsets(**options)
  request = Protocol::ListOffsetRequest.new(**options)

  send_request(request)
end

#produce(**options) ⇒ Kafka::Protocol::ProduceResponse

Produces a set of messages to the broker.

Parameters:

  • required_acks (Integer)
  • timeout (Integer)
  • messages_for_topics (Hash)

Returns:



71
72
73
74
75
# File 'lib/kafka/broker.rb', line 71

def produce(**options)
  request = Protocol::ProduceRequest.new(**options)

  send_request(request)
end

#sync_group(**options) ⇒ Object



95
96
97
98
99
# File 'lib/kafka/broker.rb', line 95

def sync_group(**options)
  request = Protocol::SyncGroupRequest.new(**options)

  send_request(request)
end

#to_sString

Returns:

  • (String)


23
24
25
# File 'lib/kafka/broker.rb', line 23

def to_s
  "#{@host}:#{@port} (node_id=#{@node_id.inspect})"
end

#txn_offset_commit(**options) ⇒ Object



191
192
193
194
195
# File 'lib/kafka/broker.rb', line 191

def txn_offset_commit(**options)
  request = Protocol::TxnOffsetCommitRequest.new(**options)

  send_request(request)
end