7H)com/dexels/kafka/impl/KafkaTopicPublisherjava/lang/Object-com/dexels/pubsub/rx2/api/PersistentPublisher(com/dexels/pubsub/rx2/api/TopicPublisherloggerLorg/slf4j/Logger;producer1Lorg/apache/kafka/clients/producer/KafkaProducer; SignatureGLorg/apache/kafka/clients/producer/KafkaProducer; adminClient,Lorg/apache/kafka/clients/admin/AdminClient; partitionsLjava/lang/Integer;replicationFactorLjava/lang/Short;detectedTopicsLjava/util/Set;#Ljava/util/Set;()VCode org/slf4j/LoggerFactory   getLogger%(Ljava/lang/Class;)Lorg/slf4j/Logger; " LineNumberTableLocalVariableTable ' % ) +java/util/HashSet *' . this+Lcom/dexels/kafka/impl/KafkaTopicPublisher;activate(Ljava/util/Map;)V8(Ljava/util/Map;)VRuntimeInvisibleAnnotations1Lorg/osgi/service/component/annotations/Activate;7java/util/Properties 6':bootstrap.servers<hosts >@? java/util/Map ABget&(Ljava/lang/Object;)Ljava/lang/Object; 6D EFput8(Ljava/lang/Object;Ljava/lang/Object;)Ljava/lang/Object;HacksJallL batch.size NPOjava/lang/Integer QRvalueOf(I)Ljava/lang/Integer;Trequest.timeout.ms_W compressionYjava/lang/String[ X] ^_equals(Ljava/lang/Object;)Zacompression.typec linger.mse buffer.memoryhkey.serializerj6org/apache/kafka/common/serialization/StringSerializer lnmjava/lang/Class opgetCanonicalName()Ljava/lang/String;rvalue.serializert9org/apache/kafka/common/serialization/ByteArraySerializervretries Nx yzparseInt(Ljava/lang/String;)I }  java/lang/Short  parseShort(Ljava/lang/String;)S Q(S)Ljava/lang/Short;   java/util/UUID  randomUUID()Ljava/util/UUID; ptoString client.id java/lang/Thread  currentThread()Ljava/lang/Thread; getContextClassLoader()Ljava/lang/ClassLoader;$org/apache/kafka/clients/KafkaClient l getClassLoader setContextClassLoader(Ljava/lang/ClassLoader;)V/org/apache/kafka/clients/producer/KafkaProducer %(Ljava/util/Properties;)V  createAdminClient=(Ljava/util/Map;)Lorg/apache/kafka/clients/admin/AdminClient;     listTopics()Ljava/util/Set;  java/util/Set addAll(Ljava/util/Collection;)ZsettingsLjava/util/Map;propsLjava/util/Properties;Ljava/lang/String;partitionsStringreplicationFactorStringclientIdoriginalLjava/lang/ClassLoader;LocalVariableTypeTable5Ljava/util/Map; StackMapTablejava/lang/ClassLoaderjava/lang/Throwable%()Ljava/util/Set; *org/apache/kafka/clients/admin/AdminClient 3()Lorg/apache/kafka/clients/admin/ListTopicsResult; /org/apache/kafka/clients/admin/ListTopicsResult names'()Lorg/apache/kafka/common/KafkaFuture; #org/apache/kafka/common/KafkaFuture A()Ljava/lang/Object;*Kafka topic publisher detected: {} topics. size()I org/slf4j/Logger info'(Ljava/lang/String;Ljava/lang/Object;)V  interrupt*Error listing kafka topics for warm cache. error*(Ljava/lang/String;Ljava/lang/Throwable;)VError listing topics: java/util/Collections emptySetjava/lang/InterruptedException'java/util/concurrent/ExecutionExceptiondetectede Ljava/lang/InterruptedException;)Ljava/util/concurrent/ExecutionException; streamTopics()Lio/reactivex/Flowable;-()Lio/reactivex/Flowable;0org/apache/kafka/clients/admin/ListTopicsOptions '     listInternal5(Z)Lorg/apache/kafka/clients/admin/ListTopicsOptions;  e(Lorg/apache/kafka/clients/admin/ListTopicsOptions;)Lorg/apache/kafka/clients/admin/ListTopicsResult; io/reactivex/Single  fromFuture4(Ljava/util/concurrent/Future;)Lio/reactivex/Single;   toFlowable apply#()Lio/reactivex/functions/Function;  io/reactivex/Flowable !"flatMap:(Lio/reactivex/functions/Function;)Lio/reactivex/Flowable;%Lorg/apache/kafka/common/KafkaFuture;JLorg/apache/kafka/common/KafkaFuture;>;c(Ljava/util/Map;)Lorg/apache/kafka/clients/admin/AdminClient;'java/util/HashMap &' >D + ,create adminSettings deactivate3Lorg/osgi/service/component/annotations/Deactivate; 1 2closeflush 5 3publisherForTopic@(Ljava/lang/String;)Lcom/dexels/pubsub/rx2/api/MessagePublisher; DeprecatedRuntimeVisibleAnnotationsLjava/lang/Deprecated;<java/lang/RuntimeException>java/lang/StringBuilder@Can not publish to topic: =B %C(Ljava/lang/String;)V =E FGappend-(Ljava/lang/String;)Ljava/lang/StringBuilder;I%, my producer is down or unconfigured = ;BM+com/dexels/kafka/impl/KafkaTopicPublisher$1 LO %P@(Lcom/dexels/kafka/impl/KafkaTopicPublisher;Ljava/lang/String;)Vtopicpublish)(Ljava/lang/String;Ljava/lang/String;[B)VU0org/apache/kafka/clients/producer/ProducerRecord TW %X9(Ljava/lang/String;Ljava/lang/Object;Ljava/lang/Object;)V Z [\sendQ(Lorg/apache/kafka/clients/producer/ProducerRecord;)Ljava/util/concurrent/Future;keyvalue[Bc(Ljava/lang/String;Ljava/lang/String;[BLjava/util/function/Consumer;Ljava/util/function/Consumer;)V(Ljava/lang/String;Ljava/lang/String;[BLjava/util/function/Consumer;Ljava/util/function/Consumer;)Vc de onCompletionh(Ljava/util/function/Consumer;Ljava/util/function/Consumer;)Lorg/apache/kafka/clients/producer/Callback; g [h}(Lorg/apache/kafka/clients/producer/ProducerRecord;Lorg/apache/kafka/clients/producer/Callback;)Ljava/util/concurrent/Future; onSuccessLjava/util/function/Consumer;onFail1Ljava/util/function/Consumer;4Ljava/util/function/Consumer; oqpjava/util/Optional rsempty()Ljava/util/Optional; u ,v=(Ljava/lang/String;Ljava/util/Optional;Ljava/util/Optional;)Vg(Ljava/lang/String;Ljava/util/Optional;Ljava/util/Optional;)V y z_contains|*Topic: {} already exists, ignoring create.~'org/apache/kafka/clients/admin/NewTopic o BorElse N intValue ()Ljava/util/function/Function; o map3(Ljava/util/function/Function;)Ljava/util/Optional;   shortValue()S } %(Ljava/lang/String;IS)V java/util/Arrays asList%([Ljava/lang/Object;)Ljava/util/List;   createTopicsK(Ljava/util/Collection;)Lorg/apache/kafka/clients/admin/CreateTopicsResult; 1org/apache/kafka/clients/admin/CreateTopicsResult J  _add)Topic: {} already exists, ignoring createError creating topic: Error: 3org/apache/kafka/common/errors/TopicExistsExceptionLjava/util/Optional;partitionCount newTopics)Lorg/apache/kafka/clients/admin/NewTopic;ctr3Lorg/apache/kafka/clients/admin/CreateTopicsResult;5Lorg/apache/kafka/common/errors/TopicExistsException;e1)Ljava/util/Optional;delete  (Ljava/util/List;)V'(Ljava/util/List;)VDeleting topics: {}   deleteTopicsK(Ljava/util/Collection;)Lorg/apache/kafka/clients/admin/DeleteTopicsResult; 1org/apache/kafka/clients/admin/DeleteTopicsResult   removeAllError deleting topics: = F-(Ljava/lang/Object;)Ljava/lang/StringBuilder;topicsLjava/util/List;3Lorg/apache/kafka/clients/admin/DeleteTopicsResult;$Ljava/util/List;java/util/List deleteGroups,(Ljava/util/List;)Lio/reactivex/Completable;@(Ljava/util/List;)Lio/reactivex/Completable;  deleteConsumerGroupsS(Ljava/util/Collection;)Lorg/apache/kafka/clients/admin/DeleteConsumerGroupsResult; 9org/apache/kafka/clients/admin/DeleteConsumerGroupsResult io/reactivex/Completable 9(Ljava/util/concurrent/Future;)Lio/reactivex/Completable;groupsdeleteCompletableDeleting topics: {} result: {}  X run\(Lcom/dexels/kafka/impl/KafkaTopicPublisher;Ljava/util/List;)Lio/reactivex/functions/Action;   doOnComplete;(Lio/reactivex/functions/Action;)Lio/reactivex/Completable; 1(Ljava/util/List;)Lio/reactivex/functions/Action;7Lorg/apache/kafka/common/KafkaFuture;backpressurePublisher7(Ljava/util/Optional;I)Lorg/reactivestreams/Subscriber;v(Ljava/util/Optional;I)Lorg/reactivestreams/Subscriber;+com/dexels/kafka/impl/KafkaTopicPublisher$2  %C(Lcom/dexels/kafka/impl/KafkaTopicPublisher;ILjava/util/Optional;)V defaultTopic maxInFlightI(Ljava/util/Optional;listConsumerGroups()Lio/reactivex/Observable;/()Lio/reactivex/Observable;  ;()Lorg/apache/kafka/clients/admin/ListConsumerGroupsResult; 7org/apache/kafka/clients/admin/ListConsumerGroupsResult valid    flatMapObservable<(Lio/reactivex/functions/Function;)Lio/reactivex/Observable; io/reactivex/Observable  describeConsumerGroups)(Ljava/util/List;)Lio/reactivex/Flowable;t(Ljava/util/List;)Lio/reactivex/Flowable;>;  U(Ljava/util/Collection;)Lorg/apache/kafka/clients/admin/DescribeConsumerGroupsResult; ;org/apache/kafka/clients/admin/DescribeConsumerGroupsResult describedGroups()Ljava/util/Map; > !"values()Ljava/util/Collection; $ %& fromIterable-(Ljava/lang/Iterable;)Lio/reactivex/Flowable; ) *" flatMapSingle, -N(Lcom/dexels/kafka/impl/KafkaTopicPublisher;)Lio/reactivex/functions/Function; / "parseGroupDescriptionJ(Lorg/apache/kafka/clients/admin/ConsumerGroupDescription;)Ljava/util/Map;p(Lorg/apache/kafka/clients/admin/ConsumerGroupDescription;)Ljava/util/Map;4groupId 6877org/apache/kafka/clients/admin/ConsumerGroupDescription 4p:state 6< :=.()Lorg/apache/kafka/common/ConsumerGroupState; ?@*org/apache/kafka/common/ConsumerGroupStateBhost 6D EF coordinator ()Lorg/apache/kafka/common/Node; HJIorg/apache/kafka/common/Node BpL memberCount =' 6O P"members RSjava/util/Collection =U FV(I)Ljava/lang/StringBuilder; X YZunmodifiableMap (Ljava/util/Map;)Ljava/util/Map;desc9Lorg/apache/kafka/clients/admin/ConsumerGroupDescription;result5Ljava/util/Map;consumerGroupOffsets)(Ljava/lang/String;)Lio/reactivex/Single;^(Ljava/lang/String;)Lio/reactivex/Single;>; c delistConsumerGroupOffsetsS(Ljava/lang/String;)Lorg/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult; gih=org/apache/kafka/clients/admin/ListConsumerGroupOffsetsResult jpartitionsToOffsetAndMetadata , m n8(Lio/reactivex/functions/Function;)Lio/reactivex/Single; describeTopicq4org/apache/kafka/clients/admin/DescribeTopicsOptions p' t uvdescribeTopics(Ljava/util/Collection;Lorg/apache/kafka/clients/admin/DescribeTopicsOptions;)Lorg/apache/kafka/clients/admin/DescribeTopicsResult; xy3org/apache/kafka/clients/admin/DescribeTopicsResult{{} descriptionTLjava/util/Map;lambda$00(Ljava/util/Set;)Lorg/reactivestreams/Publisher; Exceptionsjava/lang/Exceptionlambda$1(Ljava/util/function/Consumer;Ljava/util/function/Consumer;Lorg/apache/kafka/clients/producer/RecordMetadata;Ljava/lang/Exception;)V java/util/function/Consumer accept(Ljava/lang/Object;)VError producing: record2Lorg/apache/kafka/clients/producer/RecordMetadata; exceptionLjava/lang/Exception;lambda$2&(Ljava/lang/Integer;)Ljava/lang/Short; Nlambda$3lambda$4Deleted topics: {} complete!lambda$9 > entrySet  stream()Ljava/util/stream/Stream;  J(Lcom/dexels/kafka/impl/KafkaTopicPublisher;)Ljava/util/function/Function; java/util/stream/Stream 8(Ljava/util/function/Function;)Ljava/util/stream/Stream;   java/util/stream/Collectors toMapX(Ljava/util/function/Function;Ljava/util/function/Function;)Ljava/util/stream/Collector;  collect0(Ljava/util/stream/Collector;)Ljava/lang/Object;itemnLjava/util/Map; lambda$10I(Ljava/util/Map$Entry;)Lcom/dexels/kafka/impl/KafkaTopicPublisher$1Tuple;0com/dexels/kafka/impl/KafkaTopicPublisher$1Tuple java/util/Map$Entry getKey&org/apache/kafka/common/TopicPartition  Qp X Q&(Ljava/lang/Object;)Ljava/lang/String;-   partition  getValue3org/apache/kafka/clients/consumer/OffsetAndMetadata  offset()J  %A(Lcom/dexels/kafka/impl/KafkaTopicPublisher;Ljava/lang/String;J)VLjava/util/Map$Entry;tLjava/util/Map$Entry; lambda$11F(Lcom/dexels/kafka/impl/KafkaTopicPublisher$1Tuple;)Ljava/lang/String;  name2Lcom/dexels/kafka/impl/KafkaTopicPublisher$1Tuple; lambda$12D(Lcom/dexels/kafka/impl/KafkaTopicPublisher$1Tuple;)Ljava/lang/Long;  J java/lang/Long Q(J)Ljava/lang/Long; SourceFileKafkaTopicPublisher.java2Lorg/osgi/service/component/annotations/Component;#navajo.resource.kafkatopicpublisherconfigurationPolicy6XXX>6XXX>6XXXX>6XXXXX! >6XXXXX  K*϶ԶL!+M+L!+L!+%&%;#& pq$r&s't-u;v<wGy$*K/0'<   fT  p&*Y  L+# $&/0# $ %&Y(M,9+;=)W,*#$ /0--.4/e*( *(0*(**|#$ /036*(4# $ /06789:u0*(";Y=Y?A+DHDJKLY*+N#&$0/00Q&RS_*(TY+,-VYW# $*/0Q]^_R` a*(TY+,-VbfW# $>/0Q]^_ijkjilkm,CD *+nnt#  $ /0 Q,v w*-+x!{+}Y+-*|N,*:*}YS:W*-+WI:!+8:!=YA+DJ:!]qt]q]q#> H]fqv$\ /0QHs]^v (YXoo}P_CH*XY+S# $/0Q ]!+*+M,W*-+W4N!=YƷA+ȶJ-N!-(+(I#*  (,IJV \ $4]/0]H,J ]+] U*+Զظ۰#$/0  3*+M!+,,N-*++# )2$*3/03 *J#3J [ Y*+#$  /0     T*  #_ `a_$ /0 w%*+#'(*+.#efg$e$%/0% %01 2]&Y(M,3+5)W,9+;>)W,A+CG)W,K=YM+NQTJ)W,W#klm'n8oXp$ ]/0][\U] U]^_` a\*+bf*kl#|} |~|$/04oCJ*XY+SpYrsw>M!z,MM!,/2/<##/3<=I$4J/0JQ# |3=  # |} rI  ~A*##$   B-#*<*-2:! ++,:!(/2#. !$(/4A$*BB 4  RM 2*#$ 5 *-+#$  /0 - !*#$Zv,+*># ~$,/0, ,CY*=Y+÷AƶD+ȶTJ+ζз԰#~$C/0C C /*۰#$  2*#$ 4seZ    !"#%&)*+./034589:=>?*L@ACE>F GL