ruby-kafka gem中connection_builder是什么?应如何正确传值?
你遇到的报错是因为对Kafka::BrokerPool的connection_builder参数理解有误,该参数不接受broker地址字符串,而是要求传入可调用的连接构造对象,你传入的字符串没有hostname方法,触发了NoMethodError。
connection_builder的定义
connection_builder是ruby-kafka用来为每个broker节点构造TCP连接的可调用对象,BrokerPool会在需要和新的broker节点建立连接时,自动调用该对象,传入broker的host、port、节点id三个参数,拿到可用的Connection实例。
合法传值要求
- 必须是响应
call方法的对象,通常用Lambda/Proc实现 call方法需要接收三个入参:host(字符串,broker主机地址)、port(整数,broker端口)、broker_id(整数,broker节点id)call方法的返回值必须是初始化完成的Kafka::Connection实例
修复方案
场景1:无自定义连接需求(99%的常规场景)
普通场景下你完全不需要手动构造BrokerPool实例,直接让Kafka::Cluster生成默认的BrokerPool即可,删掉手动构造bp的代码就行:
cluster = Kafka::Cluster.new( seed_brokers: ["localhost:9092"], logger: Rails.logger, resolve_seed_brokers: false ) operation = Kafka::FetchOperation.new( cluster: cluster, logger: Rails.logger, min_bytes: 1, max_wait_time: 10 ) operation.fetch_from_partition("test-drive", 6, offset: :latest, max_bytes: 100000) operation.execute
场景2:需要自定义连接配置(比如调整超时、加SSL)
如果确实需要修改连接层的默认配置,再手动构造connection_builder和BrokerPool:
# 自定义连接构造逻辑 connection_builder = lambda do |host, port, broker_id| Kafka::Connection.new( host: host, port: port, node_id: broker_id, logger: Rails.logger, connect_timeout: 15, # 自定义连接超时 socket_timeout: 30, # 自定义请求超时 # 如有SSL需求可在这里加ssl相关配置参数 ) end bp = Kafka::BrokerPool.new( connection_builder: connection_builder, logger: Rails.logger ) cluster = Kafka::Cluster.new( seed_brokers: ["localhost:9092"], logger: Rails.logger, resolve_seed_brokers: false, broker_pool: bp ) operation = Kafka::FetchOperation.new( cluster: cluster, logger: Rails.logger, min_bytes: 1, max_wait_time: 10 ) operation.fetch_from_partition("test-drive", 6, offset: :latest, max_bytes: 100000) operation.execute
注意事项
- broker地址只需要在
seed_brokers参数中配置即可,不需要重复传给BrokerPool - 默认的connection_builder已经覆盖了绝大多数通用场景,没有特殊需求不需要手动配置该参数
内容的提问来源于stack exchange,提问作者Spoorthi kg
相关产品推荐
相关产品推荐

