【发布时间】:2016-08-30 14:15:31
【问题描述】:
confluence 文档展示了如何获取存储在 kafka 中的消费者偏移量,如下所示:https://cwiki.apache.org/confluence/display/KAFKA/Committing+and+fetching+consumer+offsets+in+Kafka
似乎分配了一个代理作为偏移管理器,所有偏移获取和提交都在这个代理上完成。但是如果这个经纪人倒闭了怎么办?
Broker offsetManager = metadataResponse.coordinator();
// if the coordinator is different, from the above channel's host then reconnect
channel.disconnect();
channel = new BlockingChannel(offsetManager.host(), offsetManager.port(),
BlockingChannel.UseDefaultBufferSize(),
BlockingChannel.UseDefaultBufferSize(),
5000 /* read timeout in millis */);
channel.connect();
【问题讨论】:
标签: apache-kafka