kafka消費者(官方示例)



import java.util.Arrays;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

public class CustomNewConsumer {

	public static void main(String[] args) {

		Properties props = new Properties();
		// 定義kakfa 服務的地址,不需要將所有broker指定上 
		props.put("bootstrap.servers", "hadoop01:9092");
		// 制定consumer group 
		props.put("group.id", "test");
		// 是否自動確認offset 
		props.put("enable.auto.commit", "true");
		// 自動確認offset的時間間隔 
		props.put("auto.commit.interval.ms", "1000");
		// key的序列化類
		props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
		// value的序列化類 
		props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
		// 定義consumer 
		KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
		
		// 消費者訂閱的topic, 可同時訂閱多個 
		consumer.subscribe(Arrays.asList("first", "second","third"));

		while (true) {
			// 讀取數據,讀取超時時間為100ms 
			ConsumerRecords<String, String> records = consumer.poll(100);
			
			for (ConsumerRecord<String, String> record : records)
				System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
		}
	}
}


免責聲明!

本站轉載的文章為個人學習借鑒使用,本站對版權不負任何法律責任。如果侵犯了您的隱私權益,請聯系本站郵箱yoyou2525@163.com刪除。



 
粵ICP備18138465號   © 2018-2025 CODEPRJ.COM