package com.dexels.kafka.impl; import java.io.FileOutputStream; import java.io.IOException; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Properties; import java.util.Set; import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import java.util.function.BiFunction; import java.util.function.Function; import java.util.stream.Collectors; import org.apache.kafka.clients.consumer.CommitFailedException; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.InterruptException; import org.apache.kafka.common.errors.WakeupException; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.osgi.service.component.annotations.Activate; import org.osgi.service.component.annotations.Component; import org.osgi.service.component.annotations.ConfigurationPolicy; import org.osgi.service.component.annotations.Deactivate; import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.reactivestreams.Subscription; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.dexels.kafka.api.OffsetQuery; import com.dexels.pubsub.rx2.api.PersistentSubscriber; import com.dexels.pubsub.rx2.api.PubSubMessage; import com.dexels.pubsub.rx2.api.TopicSubscriber; @Component(name = "navajo.resource.kafkatopicsubscriber", configurationPolicy = ConfigurationPolicy.REQUIRE, immediate = true, property={"event.topics=navajo/shutdown"}) public class KafkaTopicSubscriber implements TopicSubscriber,PersistentSubscriber, OffsetQuery { private final AtomicBoolean doCommit = new AtomicBoolean(false); private final AtomicBoolean deactived = new AtomicBoolean(false); private Properties defaultProperties; private long pollTimeout; private Properties configurationProperties; private final Set> consumers = new HashSet<>(); private final Map,Integer> lastPollCount = new HashMap<>(); private final static Logger logger = LoggerFactory.getLogger(KafkaTopicSubscriber.class); private final Map,Long> lastPoll = new HashMap<>(); private final Map> offsets = new HashMap<>(); @Activate public void activate(Map settings) { String hosts = (String) settings.get("hosts"); defaultProperties = new Properties(); defaultProperties.put("bootstrap.servers", hosts); defaultProperties.put("enable.auto.commit", false); defaultProperties.put("max.poll.records", "1"); defaultProperties.put("max.poll.interval.ms", "300000"); defaultProperties.put("auto.offset.reset", "earliest"); defaultProperties.put("request.timeout.ms", "30000"); defaultProperties.put("session.timeout.ms", "20000"); defaultProperties.put("fetch.max.wait.ms", "20000"); defaultProperties.put("key.deserializer", StringDeserializer.class.getCanonicalName()); defaultProperties.put("value.deserializer", ByteArrayDeserializer.class.getName()); this.pollTimeout = Integer.parseInt((String) settings.get("wait"));; configurationProperties = new Properties(); configurationProperties.putAll(defaultProperties); configurationProperties.put("max.poll.records", settings.get("max")); } @Deactivate public void deactivate() { deactived.set(true); for (KafkaConsumer kafkaConsumer : consumers) { closeConsumer(kafkaConsumer); } consumers.clear(); } public void markOffsets(final KafkaConsumer cons) { Map> topics = cons.listTopics(); topics.entrySet().forEach(e->{ e.getValue().forEach(p->{ }); }); } private ConsumerRecords pollConsumer(final KafkaConsumer cons, boolean allowCommit) throws InterruptException, CommitFailedException { if(allowCommit && doCommit.get()) { cons.commitSync(); doCommit.set(false); } ConsumerRecords poll = cons.poll(Duration.ofMillis(pollTimeout)); lastPollCount.put(cons, poll.count()); long now = System.currentTimeMillis(); Long previous = this.lastPoll.get(cons); // cons.subscription() if(previous!=null && poll.count()!=0) { long l = now-previous; if(l>1000) { logger.info("Time between polls: "+l+" topic: <[not supplied]> size: "+poll.count()); } } lastPoll.put(cons, now); return poll; } public Publisher> subscribe(List topics,String consumerGroup,Optional clientId, boolean fromBeginning, Runnable onPartitionsAssigned) { return subscribe(topics, consumerGroup, Optional.empty(), Optional.empty(), clientId, fromBeginning, onPartitionsAssigned); } @Override public String encodeTag(Map> tag) { return tag.entrySet() .stream() .map(e->e.getKey()+"|"+e.getValue().entrySet().stream().map(f->""+f.getKey()+":"+f.getValue()).collect(Collectors.joining(",")) ).collect(Collectors.joining(";")); } public String encodeTopicTag(Map tagMap) { return encodeTopicTag(partition->tagMap.get(partition), new ArrayList<>(tagMap.keySet())); } public String encodeTopicTag(Function tag, List partitions) { return partitions.stream().map(e->e+":"+tag.apply(e)).collect(Collectors.joining(",")); // .entrySet().stream().map(f->""+f.getKey()+":"+f.getValue()).collect(Collectors.joining(","))).collect(Collectors.joining(";")); } public Function decodeTopicTag(String tag) { String partitionPairs[] = tag.split(","); Map pairMap = new HashMap<>(); for (String pair : partitionPairs) { pairMap.put(Integer.parseInt(pair.split(":")[0]), Long.parseLong(pair.split(":")[1])); } logger.info("Decode Topic Tag: "+pairMap); return i->pairMap.get(i); } @Override public BiFunction decodeTag(String tag) { // temptopic-345|0:1,1:1;temptopic-123|0:1,1:1 Map> result = new HashMap<>(); String[] topics = tag.split(";"); for (String tp : topics) { String parts[] = tp.split("\\|"); String topic = parts[0]; result.put(topic, decodeTopicTag(parts[1])); } boolean isSingle = tag.indexOf(';')==-1; if(isSingle) { return (topic,pt)->result.entrySet().stream().findFirst().get().getValue().apply(pt); } return (topic,partition)->result.get(topic).apply(partition); } @Override public Publisher> subscribeSingleRange(String topic, String consumerGroup, String fromTag, String toTag) { final Function decodedFrom = decodeTopicTag(fromTag); final Function decodedTo = decodeTopicTag(toTag); return subscribe(Arrays.asList(new String[]{topic}), consumerGroup, Optional.of((top,part)->decodedFrom.apply(part)), Optional.of((top,part)->decodedTo.apply(part)), Optional.of("clientid-"+UUID.randomUUID().toString()), false, ()->{}); } @Override public Publisher> subscribe(List topics, String consumerGroup, Optional> fromTag, Optional> toTag, Optional clientId, boolean fromBeginning, Runnable onPartitionsAssigned, boolean commitBeforeEachPoll) { return new Publisher>() { private Subscription subscription; AtomicLong requests = new AtomicLong(); AtomicBoolean isRunning = new AtomicBoolean(false); @Override public void subscribe(Subscriber> sub) { final KafkaConsumer cons = createConsumer(consumerGroup, clientId,fromBeginning,()->{}); logger.info("Subscribed to topic: {} with groupId: {} and clientId: {}",topics,consumerGroup,clientId.orElse("")); Map> initialActivePartitions = allPartitions(topics,cons); try { cons.subscribe(topics,new KafkaConsumerRebalanceListener(cons, consumerGroup, fromBeginning, onPartitionsAssigned, fromTag)); } catch (WakeupException e) { if (!deactived.get()) throw e; } final BiConsumer commitConsumer = (topicpartition,offset)->{ OffsetAndMetadata oam = new OffsetAndMetadata(offset); Map map = new HashMap<>(); map.put(topicpartition, oam); cons.commitAsync(map, (result,exception)->{}); }; this.subscription = new Subscription() { Map> activePartitions = initialActivePartitions; AtomicBoolean isCompleted = new AtomicBoolean(false); @Override public void request(long r) { Throwable detectedError = null; requests.addAndGet(r); try { while(requests.get()>0 && isRunning.get() && detectedError==null && isCompleted.get()==false) { requests.decrementAndGet(); ConsumerRecords result; try { result = pollConsumer(cons,commitBeforeEachPoll); doCommit.set(true); // logger.info("Consumer records: "+result.count()); List res = new ArrayList<>(result.count()); for (ConsumerRecord e : result) { if(!isRunning.get()) { break; } final KafkaMessage kafkaMsg = new KafkaMessage(e,commitConsumer); final Set set = activePartitions.get(kafkaMsg.topic().get()); if(set!=null) { if(set.contains(kafkaMsg.partition().get())) { res.add(kafkaMsg); } } activePartitions = continueAfter(kafkaMsg, toTag, activePartitions); if (activePartitions.isEmpty()) { logger.info("No more active partitions for groupId: {}", consumerGroup); isRunning.set(false); } Map topicOffsets = offsets.get(e.topic()); if(topicOffsets==null) { topicOffsets = new HashMap<>(); offsets.put(e.topic(), topicOffsets); } topicOffsets.put(e.partition(), e.offset()); } sub.onNext(res); } catch (InterruptException e1) { detectedError = e1; logger.error("Interrupted at groupId {} : ", consumerGroup, e1); isRunning.set(false); } catch(CommitFailedException cfe) { sub.onError(cfe); logger.error("Commit issue at groupId: {} downstream should re-subscribe.", consumerGroup); isRunning.set(false); } } } catch (Throwable e) { logger.error("Error in poll for groupId (): ", consumerGroup, e); sub.onError(e); isRunning.set(false); isCompleted.set(true); return; } if(!isRunning.get()) { if(!isCompleted.get()) { sub.onComplete(); subscription.cancel(); isCompleted.set(true); } return; } } @Override public void cancel() { logger.info("Cancelling subscription for groupId: {}", consumerGroup); isRunning.set(false); } }; isRunning.set(true); sub.onSubscribe(subscription); } }; } protected Map> allPartitions(List topics, KafkaConsumer cons) { Map> result = new HashMap<>(); topics.forEach(topic->{ Set partitions = cons.partitionsFor(topic) .stream() .map(p->p.partition()) .collect(Collectors.toSet()); result.put(topic, partitions); }); return result; } private Map> continueAfter(KafkaMessage msg, Optional> toTag, Map> activePartitions) { if(!toTag.isPresent()) { return activePartitions; } if(!msg.topic().isPresent() || !msg.partition().isPresent() || !msg.offset().isPresent()) { throw new RuntimeException("Missing topic, partition or offset in message"); } Long offset = toTag.get().apply(msg.topic().get(), msg.partition().get()); if(offset==null) { logger.info("Weird: No offset found for topic: {} and partition: {}",msg.topic().get(), msg.partition().get()); return activePartitions; } if(msg.offset().get()+1 >= offset) { logger.info("Message past target offset: {} : topic: {} partition: {} offset: {}",offset, msg.topic().get(), msg.partition().get(), msg.offset().get()); // partitionCompletedConsumer.accept(msg.topic().get(), msg.partition().get()); return partitionDone(msg.topic().get(), msg.partition().get(), activePartitions); } return activePartitions; } private Map> partitionDone(String topic, int partition, Map> activePartitions) { Set activeNow = new HashSet<>(activePartitions.get(topic)); if(!activeNow.contains(partition)) { // already removed? return activePartitions; // throw new RuntimeException("Deactivating inactive partition: "+activePartitions+" topic: "+topic+" partitions: "+partition); } Map> result = new HashMap<>(activePartitions); activeNow.remove(partition); if(activeNow.isEmpty()) { result.remove(topic); return result; } result.put(topic, activeNow); return result; } private synchronized KafkaConsumer createConsumer(String groupId, Optional clientId, boolean withHistory, Runnable onPartitionsAssigned) { KafkaConsumer consumer; logger.info("Creating a consumer groupid: {} clientid: {} with history: {}",groupId,clientId,withHistory); final ClassLoader original = Thread.currentThread().getContextClassLoader(); try { Properties copy = new Properties(); copy.putAll(configurationProperties); copy.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); if(clientId.isPresent()) { copy.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId.get()); } if(withHistory) { copy.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); } else { copy.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); } Thread.currentThread().setContextClassLoader(KafkaConsumer.class.getClassLoader()); consumer = new KafkaConsumer<>(copy); consumers.add(consumer); } finally { Thread.currentThread().setContextClassLoader(original); } return consumer; } private synchronized void closeConsumer(KafkaConsumer consumer) { logger.info("Closing consumer on thread: {}", Thread.currentThread().getName()); consumer = null; } private AtomicLong written = new AtomicLong(); private static long started = System.currentTimeMillis(); protected void writeToFile(FileOutputStream fo, byte[] a) { try { long w = written.addAndGet(a.length); fo.write(a); if(w % 1000==0) { logger.debug("Written: "+written); long elapsed = System.currentTimeMillis() - started; float rate = (float)w / (float)elapsed / 1024f / 1024f * 1000; logger.debug("rate: {} MB/s", rate); } } catch (IOException e) { e.printStackTrace(); } } private final class KafkaConsumerRebalanceListener implements ConsumerRebalanceListener { private KafkaConsumer consumer; private String groupId; private boolean withHistory; private Runnable onPartitionsAssigned; private final Optional> fromTag; public KafkaConsumerRebalanceListener(KafkaConsumer consumer, String groupId, boolean withHistory, Runnable onPartitionsAssigned,Optional> fromTag) { this.consumer = consumer; this.groupId = groupId; this.withHistory = withHistory; this.onPartitionsAssigned = onPartitionsAssigned; this.fromTag = fromTag; } @Override public void onPartitionsRevoked(Collection c) { for (TopicPartition tp : c) { Long last = lastPoll.get(consumer); if(last!=null) { logger.info("Partitions revoked: {} millis after last poll!",(System.currentTimeMillis()- last)); } logger.info("Partition revoked for topic: {} and partition: {}",tp.topic(),tp.partition()); } } @Override public void onPartitionsAssigned(Collection c) { logger.info("Partitions assigned!"); if(fromTag.isPresent()) { BiFunction from = fromTag.get(); for (TopicPartition tp : c) { Long offset = from.apply(tp.topic(), tp.partition()); logger.info("GroupId: {} Partition assigned for topic: {} and partition: {} at position: {} seeking to: {}",groupId,tp.topic(),tp.partition(),consumer.position(tp),offset); consumer.seek(tp, offset); } } else { if(!withHistory) { logger.info("Not with history"); boolean only0 = true; for (TopicPartition tp : c) { long position = consumer.position(tp); logger.info("GroupId: {} Topic: {} Partition: {} no history selected. Current pos: {}",groupId,tp.topic(),tp.partition(),position); if(position!=0) { only0 = false; } } if (only0) { logger.info("Only position 0 detected for group: {}, so assuming a new topic / group, fast forwarding as we're not interested in history",groupId); consumer.seekToEnd(c); for (TopicPartition tp : c) { logger.info("GroupId: {} Topic: {} Partition: {} fast forwarded. Current pos: {}",groupId,tp.topic(),tp.partition(), consumer.position(tp)); } consumer.commitSync(); } else { logger.info("Position offsets detected, so won't fast forward group: {}",groupId); } } } onPartitionsAssigned.run(); } } @Override public Publisher> subscribe(String topic, String consumerGroup, boolean fromBeginning) { return subscribe(Arrays.asList(new String[]{topic}), consumerGroup, Optional.empty(), fromBeginning, ()->{}); } @Override public List topics() { KafkaConsumer consumer = createConsumer( "offsetquery", Optional.of("offsetquery-"+UUID.randomUUID().toString()), false, ()->{}); List topics = consumer.listTopics() .entrySet() .stream() .map(pp->pp.getKey()) .collect(Collectors.toList()); consumer.close(); return topics; } public Map partitionOffsets(String topic) { KafkaConsumer consumer = createConsumer("offsetquery", Optional.of("offsetquery-"+UUID.randomUUID().toString()), false, ()->{}); Map offsets = offsetsForTopic(topic, consumer); consumer.close(); consumers.remove(consumer); return offsets; } private Map offsetsForTopic(String topic, KafkaConsumer consumer) { List parts = consumer.partitionsFor(topic) .stream() .map(pp->new TopicPartition(pp.topic(), pp.partition())) .collect(Collectors.toList()); Map offsets = consumer.endOffsets(parts) .entrySet() .stream() .collect( Collectors.toMap(e->e.getKey().partition(), f->f.getValue().longValue()+1) ); return offsets; } @Override public Map> offsets(List topics) { KafkaConsumer consumer = createConsumer( "offsetquery", Optional.of("offsetquery-"+UUID.randomUUID().toString()), false, ()->{}); return topics.stream() .collect(Collectors.toMap(Function.identity(), s->offsetsForTopic(s, consumer))); } }