echo "Waiting for the Kafka-Cluster to become ready..."
docker-compose exec cli cub kafka-ready -b kafka:9092 1 60 > /dev/null 2>&1 || exit 1
-docker-compose up -d kafka-ui
+docker-compose up setup
+docker-compose up -d producer peter beate
-docker-compose exec -T cli bash << 'EOF'
-echo "Creating topic with 3 partitions..."
-kafka-topics --bootstrap-server kafka:9092 --delete --if-exists --topic test
-# tag::createtopic[]
-kafka-topics --bootstrap-server kafka:9092 --create --topic test --partitions 3
-# end::createtopic[]
-kafka-topics --bootstrap-server kafka:9092 --describe --topic test
-EOF
+sleep 15
-docker-compose up -d consumer
-
-docker-compose up -d producer
+http -v post :8082/stop
sleep 10
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-
-docker-compose stop producer
-docker-compose exec -T cli bash << 'EOF'
-echo "Altering number of partitions from 3 to 7..."
-# tag::altertopic[]
-kafka-topics --bootstrap-server kafka:9092 --alter --topic test --partitions 7
-kafka-topics --bootstrap-server kafka:9092 --describe --topic test
-# end::altertopic[]
-EOF
+docker-compose kill -s 9 peter
+http -v post :8082/start
+sleep 60
-docker-compose start producer
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-sleep 1
-http -v :8081/seen
-docker-compose stop producer consumer
+docker-compose stop producer peter beate
+docker-compose logs beate
+docker-compose logs --tail=10 peter
ME_CONFIG_MONGODB_ADMINPASSWORD: training
ME_CONFIG_MONGODB_URL: mongodb://juplo:training@mongo:27017/
- kafka-ui:
- image: provectuslabs/kafka-ui:0.3.3
- ports:
- - 8080:8080
- environment:
- KAFKA_CLUSTERS_0_NAME: local
- KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
+ setup:
+ image: juplo/toolbox
+ command: >
+ bash -c "
+ kafka-topics --bootstrap-server kafka:9092 --delete --if-exists --topic test
+ kafka-topics --bootstrap-server kafka:9092 --create --topic test --partitions 2
+ "
cli:
image: juplo/toolbox
producer.bootstrap-server: kafka:9092
producer.client-id: producer
producer.topic: test
- producer.throttle-ms: 10
+ producer.throttle-ms: 500
- consumer:
+ peter:
image: juplo/endless-consumer:1.0-SNAPSHOT
ports:
- 8081:8081
environment:
consumer.bootstrap-server: kafka:9092
- consumer.client-id: consumer
+ consumer.client-id: peter
+ consumer.topic: test
+ spring.data.mongodb.uri: mongodb://juplo:training@mongo:27017
+ spring.data.mongodb.database: juplo
+
+ beate:
+ image: juplo/endless-consumer:1.0-SNAPSHOT
+ ports:
+ - 8082:8081
+ environment:
+ consumer.bootstrap-server: kafka:9092
+ consumer.client-id: beate
consumer.topic: test
spring.data.mongodb.uri: mongodb://juplo:training@mongo:27017
spring.data.mongodb.database: juplo
props.put("bootstrap.servers", bootstrapServer);
props.put("group.id", groupId);
props.put("client.id", id);
+ props.put("enable.auto.commit", false);
props.put("auto.offset.reset", autoOffsetReset);
props.put("metadata.max.age.ms", "1000");
props.put("key.deserializer", StringDeserializer.class.getName());
tp.partition(),
key);
}
- repository.save(new StatisticsDocument(tp.partition(), removed));
+ repository.save(new StatisticsDocument(tp.partition(), removed, consumer.position(tp)));
});
}
partitions.forEach(tp ->
{
log.info("{} - adding partition: {}", id, tp);
- seen.put(
- tp.partition(),
+ StatisticsDocument document =
repository
.findById(Integer.toString(tp.partition()))
- .map(document -> document.statistics)
- .orElse(new HashMap<>()));
+ .orElse(new StatisticsDocument(tp.partition()));
+ consumer.seek(tp, document.offset);
+ seen.put(tp.partition(), document.statistics);
});
}
});
seenByKey++;
byKey.put(key, seenByKey);
}
+
+ seen.forEach((partiton, statistics) -> repository.save(
+ new StatisticsDocument(
+ partiton,
+ statistics,
+ consumer.position(new TopicPartition(topic, partiton)))));
}
}
catch(WakeupException e)
{
@Id
public String id;
+ public long offset;
public Map<String, Integer> statistics;
public StatisticsDocument()
{
}
- public StatisticsDocument(Integer partition, Map<String, Integer> statistics)
+ public StatisticsDocument(Integer partition)
+ {
+ this.id = Integer.toString(partition);
+ this.statistics = new HashMap<>();
+ }
+
+ public StatisticsDocument(Integer partition, Map<String, Integer> statistics, long offset)
{
this.id = Integer.toString(partition);
this.statistics = statistics;
+ this.offset = offset;
}
}