import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;

import java.io.IOException;
import java.nio.charset.Charset;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.time.LocalDate;
import java.time.ZoneOffset;
import java.util.*;
import java.util.concurrent.ExecutionException;
import java.util.function.BiConsumer;
import java.util.stream.Collectors;

public class TopicsStats {

    public static final String START_DATE = "2018-12-01";
    public static final String END_DATE = "2019-01-10";

    private static class TopicStatistics {
        private String name;
        private LocalDate date;
        private Long count;

        public TopicStatistics(String name, LocalDate date, Long count) {
            this.name = name;
            this.date = date;
            this.count = count;
        }

        @Override
        public String toString() {
            return "TopicStatistics{" +
                    "name='" + name + '\'' +
                    ", date=" + date +
                    ", count=" + count +
                    '}';
        }

        public String toCSV() {
            return name + ";" + date + ";" + count + "\n";
        }
    }

    public static void main(String[] args) {

        String bootstrapServers = "broker0:9092";

        Properties config = new Properties();
        config.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        AdminClient adminClient = AdminClient.create(config);

        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<>(props);

        Map<LocalDate, Long> localDateLongMap = new TreeMap<>();

        List<TopicStatistics> topicStatistics = new ArrayList<>();

        Set<String> topics = new HashSet<>();
        try {
            topics = adminClient.listTopics().names().get().stream().filter(s -> s.startsWith("someprefix.") .collect(Collectors.toSet());
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        }

        for (String topic : topics) {
            System.out.println(topic);

            for (LocalDate date = LocalDate.parse(START_DATE).minusDays(1); date.isBefore(LocalDate.parse(END_DATE)); date = date.plusDays(1)) {
                Map<TopicPartition, Long> topicPartitionLongMap = new HashMap<>();
                LocalDate finalStartDate = date;
                kafkaConsumer.partitionsFor(topic).stream().map(partitionInfo -> new TopicPartition(partitionInfo.topic(), partitionInfo.partition())).forEach(topicPartition -> topicPartitionLongMap.put(topicPartition, finalStartDate.atStartOfDay().toEpochSecond(ZoneOffset.UTC) * 1000));

                localDateLongMap.put(finalStartDate,
                        kafkaConsumer.offsetsForTimes(topicPartitionLongMap).entrySet().stream()
                                .filter(topicPartitionOffsetAndTimestampEntry -> topicPartitionOffsetAndTimestampEntry.getValue() != null)
                                .mapToLong(topicPartitionOffsetAndTimestampEntry -> topicPartitionOffsetAndTimestampEntry.getValue().offset()).sum());

            }

            localDateLongMap.forEach(new BiConsumer<>() {
                boolean first = true;
                long lastValue = 0;

                @Override
                public void accept(LocalDate localDate, Long aLong) {
                    if (!first) {
                        long diff = aLong - lastValue;
                        topicStatistics.add(new TopicStatistics(topic, localDate, diff));
                    }
                    lastValue = aLong;
                    first = false;
                }
            });

        }

        List<String> lines = topicStatistics.stream().map(tp -> tp.toCSV()).collect(Collectors.toList());
        Path file = Paths.get("report.csv");
        try {
            Files.write(file, lines, Charset.forName("UTF-8"));
        } catch (IOException e) {
            e.printStackTrace();
        }

        kafkaConsumer.close();
        adminClient.close();

    }
}
