Advertisement
Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- public class kafkaconsumer {
- public static void main(String[] args) throws IOException {
- String TOPIC_NAME = "test";
- Properties props = new Properties();
- props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
- "localhost:9092");
- props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
- "org.apache.kafka.common.serialization.IntegerDeserializer");
- props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
- "org.apache.kafka.common.serialization.StringDeserializer");
- props.put(ConsumerConfig.GROUP_ID_CONFIG, "test" );
- props.put("enable.auto.commit", "false");
- props.put("auto.commit.interval.ms", "1000");
- props.put("session.timeout.ms", "30000");
- props.put("partition.assignment.strategy", "range");
- Consumer<String, String> consumer = new KafkaConsumer<> .
- (props);
- consumer.subscribe(Collections.singletonList("test"));
- ConsumerRecords<String, String> records =
- consumer.poll(100);
- System.out.println(consumer);
- System.out.println(records);
- consumer.close();
- }
- }
Advertisement
Add Comment
Please, Sign In to add comment
Advertisement