kafka consumer group总结
kafka消费者api分为high api和low api,目前上述demo是都是使用kafka high api,高级api不用关心维护消费状态信息和负载均衡,不用关心offset。
高级api的一些注意事项:
1. 如果consumer group中的consumer线程数量比partition多,那么有的线程将永远不会收到消息。
因为kafka的设计是在一个partition上是不允许并发的,所以consumer数不要大于partition数
高级api的一些注意事项:
1. 如果consumer group中的consumer线程数量比partition多,那么有的线程将永远不会收到消息。
因为kafka的设计是在一个partition上是不允许并发的,所以consumer数不要大于partition数
2,如果consumer group中的consumer线程数量比partition少,那么有的线程将会收到多个消息。并且不保证数据间的顺序性,kafka只保证在一个partition上数据是有序的,
3,增减consumer,broker,partition会导致rebalance,所以rebalance后consumer对应的partition会发生变化
4,High-level接口中获取不到数据的时候是会block的
关于consumer group(high api)的几点总结:
1,以consumer group为单位订阅 topic,每个consumer一起去消费一个topic;
2,consumer group 通过zookeeper来消费kafka集群中的消息(这个过程由zookeeper进行管理);
相对于low api自己管理offset,high api把offset的管理交给了zookeeper,但是high api并不是消费一次就在zookeeper中更新一次,而是每间隔一个(默认1000ms)时间更新一次offset,可能在重启消费者时拿到重复的消息。此外,当分区leader发生变更时也可能拿到重复的消息。因此在关闭消费者时最好等待一定时间(10s)然后再shutdown。
3,consumer group 设计的目的之一也是为了应用多线程同时去消费一个topic中的数据。
例子:
-
import kafka.consumer.ConsumerIterator;
-
import kafka.consumer.KafkaStream;
-
-
public class ConsumerTest implements Runnable {
-
private KafkaStream m_stream;
-
private int m_threadNumber;
-
-
public ConsumerTest(KafkaStream a_stream, int a_threadNumber) {
-
m_threadNumber = a_threadNumber;
-
m_stream = a_stream;
-
}
-
-
public void run() {
-
ConsumerIterator<byte[], byte[]> it = m_stream.iterator();
-
while (it.hasNext())
-
System.out.println(“Thread “ + m_threadNumber + “: “ + new String(it.next().message()));
-
System.out.println(“Shutting down Thread: “ + m_threadNumber);
-
}
-
}
-
-
//配置连接zookeeper的信息
-
private static ConsumerConfig createConsumerConfig(String a_zookeeper, String a_groupId) {
-
Properties props = new Properties();
-
props.put(“zookeeper.connect”, a_zookeeper); //zookeeper连接地址
-
props.put(“group.id”, a_groupId); //consumer group的id
-
props.put(“zookeeper.session.timeout.ms”, “400”);
-
props.put(“zookeeper.sync.time.ms”, “200”);
-
props.put(“auto.commit.interval.ms”, “1000”);
-
return new ConsumerConfig(props);
-
}
-
-
//建立一个消费者线程池
-
public void run(int a_numThreads) {
-
Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
-
topicCountMap.put(topic, new Integer(a_numThreads));
-
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap);
-
List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);
-
-
-
// now launch all the threads
-
//
-
executor = Executors.newFixedThreadPool(a_numThreads);
-
-
// now create an object to consume the messages
-
//
-
int threadNumber = 0;
-
for (final KafkaStream stream : streams) {
-
executor.submit(new ConsumerTest(stream, threadNumber));
-
threadNumber++;
-
}
-
}
-
-
//经过一段时间后关闭
-
try {
-
Thread.sleep(10000);
-
} catch (InterruptedException ie) {
-
-
}
-
example.shutdown();
-
import kafka.consumer.ConsumerIterator;
-
import kafka.consumer.KafkaStream;
-
-
public class ConsumerTest implements Runnable {
-
private KafkaStream m_stream;
-
private int m_threadNumber;
-
-
public ConsumerTest(KafkaStream a_stream, int a_threadNumber) {
-
m_threadNumber = a_threadNumber;
-
m_stream = a_stream;
-
}
-
-
public void run() {
-
ConsumerIterator<byte[], byte[]> it = m_stream.iterator();
-
while (it.hasNext())
-
System.out.println(“Thread “ + m_threadNumber + “: “ + new String(it.next().message()));
-
System.out.println(“Shutting down Thread: “ + m_threadNumber);
-
}
-
}
-
-
//配置连接zookeeper的信息
-
private static ConsumerConfig createConsumerConfig(String a_zookeeper, String a_groupId) {
-
Properties props = new Properties();
-
props.put(“zookeeper.connect”, a_zookeeper); //zookeeper连接地址
-
props.put(“group.id”, a_groupId); //consumer group的id
-
props.put(“zookeeper.session.timeout.ms”, “400”);
-
props.put(“zookeeper.sync.time.ms”, “200”);
-
props.put(“auto.commit.interval.ms”, “1000”);
-
return new ConsumerConfig(props);
-
}
-
-
//建立一个消费者线程池
-
public void run(int a_numThreads) {
-
Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
-
topicCountMap.put(topic, new Integer(a_numThreads));
-
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap);
-
List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);
-
-
-
// now launch all the threads
-
//
-
executor = Executors.newFixedThreadPool(a_numThreads);
-
-
// now create an object to consume the messages
-
//
-
int threadNumber = 0;
-
for (final KafkaStream stream : streams) {
-
executor.submit(new ConsumerTest(stream, threadNumber));
-
threadNumber++;
-
}
-
}
-
-
//经过一段时间后关闭
-
try {
-
Thread.sleep(10000);
-
} catch (InterruptedException ie) {
-
-
}
-
example.shutdown();
转载请注明:SuperIT » kafka consumer group总结