7'com/dexels/mongodb/sink/RunKafkaConnectjava/lang/Object)com/dexels/kafka/streams/base/ConnectSinkloggerLorg/slf4j/Logger;worker)Lorg/apache/kafka/connect/runtime/Worker;herder>Lorg/apache/kafka/connect/runtime/standalone/StandaloneHerder;shutdown+Ljava/util/concurrent/atomic/AtomicBoolean; stopLatch%Ljava/util/concurrent/CountDownLatch;offsetBackingStore5Lorg/apache/kafka/connect/storage/OffsetBackingStore;connectorConfigs[Ljava/util/Properties;workerPropertiesLjava/util/Properties;sinkProperties()VCode org/slf4j/LoggerFactory   getLogger%(Ljava/lang/Class;)Lorg/slf4j/Logger; " LineNumberTableLocalVariableTable(Ljava/util/Optional;Lcom/dexels/kafka/streams/base/StreamConfiguration;Lcom/dexels/kafka/streams/api/TopologyContext;Ljava/lang/String;Ljava/io/File;)V Exceptions)java/io/IOException Signature(Ljava/util/Optional;Lcom/dexels/kafka/streams/base/StreamConfiguration;Lcom/dexels/kafka/streams/api/TopologyContext;Ljava/lang/String;Ljava/io/File;)V - %/)java/util/concurrent/atomic/AtomicBoolean .1 %2(Z)V 4 6#java/util/concurrent/CountDownLatch 58 %9(I)V ;  =?>1com/dexels/kafka/streams/base/StreamConfiguration @Asink((Ljava/lang/String;)Ljava/util/Optional; CEDjava/util/Optional FG isPresent()ZI"java/lang/IllegalArgumentExceptionKjava/lang/StringBuilderMNo sink found: JO %P(Ljava/lang/String;)V JR STappend-(Ljava/lang/String;)Ljava/lang/StringBuilder; JV WXtoString()Ljava/lang/String; HO[connect- ]_^*com/dexels/kafka/streams/api/CoreOperators `agenerationalGroupT(Ljava/lang/String;Lcom/dexels/kafka/streams/api/TopologyContext;)Ljava/lang/String; Cc deget()Ljava/lang/Object;g.com/dexels/kafka/streams/xml/parser/XMLElement fi jk getChildren()Ljava/util/Vector;m3com/dexels/kafka/streams/api/sink/SinkConfiguration o pq createIndicesv(Ljava/util/List;Lcom/dexels/kafka/streams/api/TopologyContext;Lcom/dexels/kafka/streams/api/sink/SinkConfiguration;)V s tu sinkTopicO(Ljava/util/List;Lcom/dexels/kafka/streams/api/TopologyContext;)Ljava/util/Map; wyxjava/util/Collections z{emptyMap()Ljava/util/Map; l} ~{settingsjava/util/Properties -    sink.properties java/lang/Class getResourceAsStream)(Ljava/lang/String;)Ljava/io/InputStream;  load(Ljava/io/InputStream;)Vstandalone.propertiesname  put8(Ljava/lang/Object;Ljava/lang/Object;)Ljava/lang/Object;sinkidconnectionString  java/util/Map d&(Ljava/lang/Object;)Ljava/lang/Object;databasejava/lang/String?Resolved mongo sink definition to {} with generational group {} org/slf4j/Logger info9(Ljava/lang/String;Ljava/lang/Object;Ljava/lang/Object;)Vmongodb.database entrySet()Ljava/util/Set;  java/util/Set stream()Ljava/util/stream/Stream; acceptR(Lcom/dexels/kafka/streams/base/StreamConfiguration;)Ljava/util/function/Consumer; java/util/stream/Stream forEach (Ljava/util/function/Consumer;)Vtopics, apply()Ljava/util/function/Function; map8(Ljava/util/function/Function;)Ljava/util/stream/Stream; java/util/stream/Collectors toList()Ljava/util/stream/Collector; collect0(Ljava/util/stream/Collector;)Ljava/lang/Object;java/lang/Iterable join@(Ljava/lang/CharSequence;Ljava/lang/Iterable;)Ljava/lang/String;mongodb.collections generation ,com/dexels/kafka/streams/api/TopologyContext Ljava/lang/String; instanceName instancesinkName tenantLjava/util/Optional; deployment = Xbootstrap.serversjava/lang/CharSequence = X kafkaHosts  E(Ljava/lang/CharSequence;[Ljava/lang/CharSequence;)Ljava/lang/String;group.idcompression.type lz4  java/io/File mongooffset-   %#(Ljava/io/File;Ljava/lang/String;)Voffset.storage.file.filename   XgetAbsolutePathSink: {}  '(Ljava/lang/String;Ljava/lang/Object;)V Worker: {}   P setupConfigthis)Lcom/dexels/mongodb/sink/RunKafkaConnect;xconfig3Lcom/dexels/kafka/streams/base/StreamConfiguration;topologyContext.Lcom/dexels/kafka/streams/api/TopologyContext; storageFolderLjava/io/File; sinkConfigsinkId topicMappingLjava/util/Map;resolvedDatabaseoutputLocalVariableTypeTableFLjava/util/Optional;KLjava/util/Optional;5Ljava/util/Map; StackMapTable(Ljava/util/List;Lcom/dexels/kafka/streams/api/TopologyContext;)Ljava/util/Map;7java/util/HashMap 6-:java/util/ArrayList 9- =?>java/util/List @Aiterator()Ljava/util/Iterator; CEDjava/util/Iterator FenextHgenerationcollection fJ KLgetStringAttribute&(Ljava/lang/String;)Ljava/lang/String; N OPvalueOf&(Ljava/lang/Object;)Ljava/lang/String;R collection =T UVadd(Ljava/lang/Object;)ZXtopic ]Z [a topicName].Connecting topic: {} to mongodb collection: {} C` aGhasNext wc deunmodifiableMap (Ljava/util/Map;)Ljava/util/Map;Ljava/util/List;result collectionse0Lcom/dexels/kafka/streams/xml/parser/XMLElement;gencollcollBLjava/util/List;$Ljava/util/List; prq"org/apache/kafka/common/utils/Time stSYSTEM$Lorg/apache/kafka/common/utils/Time;v0Kafka Connect standalone worker initializing ... x Pz+org/apache/kafka/connect/runtime/WorkerInfo y- y} ~logAll java/lang/Thread  currentThread()Ljava/lang/Thread;  getContextClassLoader()Ljava/lang/ClassLoader;  getClass()Ljava/lang/Class;  getClassLoader  setContextClassLoader(Ljava/lang/ClassLoader;)V #org/apache/kafka/common/utils/Utils propsToStringMap'(Ljava/util/Properties;)Ljava/util/Map;;>; .  getAndSet(Z)ZKafka ConnectEmbedded stopping  stop   Kafka ConnectEmbedded stopped 5  countDownwasShuttingDownZ(Ljava/util/List;Lcom/dexels/kafka/streams/api/TopologyContext;Lcom/dexels/kafka/streams/api/sink/SinkConfiguration;)V = GisEmpty 'com/dexels/mongodb/sink/MongodbSinkTask  createClient4(Ljava/lang/String;)Lcom/mongodb/client/MongoClient; com/mongodb/client/MongoClient   getDatabase6(Ljava/lang/String;)Lcom/mongodb/client/MongoDatabase; "$# com/mongodb/client/MongoDatabase %& getCollection8(Ljava/lang/String;)Lcom/mongodb/client/MongoCollection;( sink.index f* +,getChildrenByTagName$(Ljava/lang/String;)Ljava/util/List;.unique0true2false f4 56getBooleanAttribute:(Ljava/lang/String;Ljava/lang/String;Ljava/lang/String;Z)Z8sparse: sphereVersion f< => getAttribute&(Ljava/lang/String;)Ljava/lang/Object;@partialFilterExpressionB definitionDiCreating index for collection: {} with definition: {} unique? {} sparse? {} partial: {} sphereVersion: {} FHGjava/lang/Boolean OI(Z)Ljava/lang/Boolean; K L((Ljava/lang/String;[Ljava/lang/Object;)V N OP createIndexp(Lcom/mongodb/client/MongoCollection;Ljava/lang/String;ZZLjava/lang/String;Ljava/lang/String;)Ljava/lang/String;RIndex result: {} T Uclose sinkElements5Lcom/dexels/kafka/streams/api/sink/SinkConfiguration;client Lcom/mongodb/client/MongoClient;db"Lcom/mongodb/client/MongoDatabase;collectionName$Lcom/mongodb/client/MongoCollection; indexElementsindex9Lcom/mongodb/client/MongoCollection;b"com/mongodb/client/MongoCollection(Lcom/mongodb/client/MongoCollection;Ljava/lang/String;ZZLjava/lang/String;Ljava/lang/String;)Ljava/lang/String; e fgsplit'(Ljava/lang/String;)[Ljava/lang/String; ikjjava/util/Arrays lmasList%([Ljava/lang/Object;)Ljava/util/List;oorg/bson/Document n-r%com/mongodb/client/model/IndexOptions q- qu .v*(Z)Lcom/mongodb/client/model/IndexOptions; qx 8v nz {|parse'(Ljava/lang/String;)Lorg/bson/Document; q~ @D(Lorg/bson/conversions/Bson;)Lcom/mongodb/client/model/IndexOptions;2Invalid document given as partial filter index! {}  warn  parseInt(Ljava/lang/String;)I q :<(Ljava/lang/Integer;)Lcom/mongodb/client/model/IndexOptions;: = d(I)Ljava/lang/Object;-?\d+  matches(Ljava/lang/String;)Z n 8(Ljava/lang/String;Ljava/lang/Object;)Ljava/lang/Object;_   replaceAll8(Ljava/lang/String;Ljava/lang/String;)Ljava/lang/String; a  listIndexes*()Lcom/mongodb/client/ListIndexesIterable; ?&com/mongodb/client/ListIndexesIterable n  Vequals3Created index on collection: {} with definition: {} a  getNamespace()Lcom/mongodb/MongoNamespace; com/mongodb/MongoNamespace XgetCollectionName a OV(Lorg/bson/conversions/Bson;Lcom/mongodb/client/model/IndexOptions;)Ljava/lang/String;Created index result: {}partialpartdLorg/bson/Document; indexOptions'Lcom/mongodb/client/model/IndexOptions; partialIndextLjava/lang/Throwable;ptprts directionI searchTermdoclambda$0K(Lcom/dexels/kafka/streams/base/StreamConfiguration;Ljava/util/Map$Entry;)V =  adminClient.()Lorg/apache/kafka/clients/admin/AdminClient; java/util/Map$Entry egetKey ] topicPartitionCount()I ] topicReplicationCount )com/dexels/kafka/streams/tools/KafkaUtils ensureExistsSyncC(Lorg/apache/kafka/clients/admin/AdminClient;Ljava/lang/String;II)VLjava/util/Map$Entry;;Ljava/util/Map$Entry;lambda$1)(Ljava/util/Map$Entry;)Ljava/lang/String;lambda$2  egetValue SourceFileRunKafkaConnect.javaBootstrapMethods "java/lang/invoke/LambdaMetafactory  metafactory(Ljava/lang/invoke/MethodHandles$Lookup;Ljava/lang/String;Ljava/lang/invoke/MethodType;Ljava/lang/invoke/MethodType;Ljava/lang/invoke/MethodHandle;Ljava/lang/invoke/MethodType;)Ljava/lang/invoke/CallSite;(Ljava/lang/Object;)V  (Ljava/util/Map$Entry;)V     InnerClasses %java/lang/invoke/MethodHandles$Lookup java/lang/invoke/MethodHandlesLookupEntry/org/apache/kafka/connect/runtime/Herder$Created'org/apache/kafka/connect/runtime/HerderCreated NestMembers!     ) !#4$%&'(*+. *,*.Y03*5Y7:,<:BHYJYLNQUYJYZNQU-\::+B+bfh-bln+B*+bfh-rv: bl|: *Y*Y***W*W* W -\: ! * W ,* ѹ۸ݶW* ѹ۸ݶW*-W*-W*W-B*-bW*,W*Y,SW*W* W YJY N QU: * W!*!**-#(A9:C$D,EEG\H`JgK}NPQRSTUVWXY%Z1[H^x_abcdegjk l.nKo[pjqyrs$ !"#$%&'()$^*\&``"+,- ~- k. K7/) 0*#1$^*2,3 ~3 48EC= C7BQtu*5 6Y8N9Y;:+<:sBf:GI:JY,MNQU QI:SWWI,Y: !\ - ^W_-b#. xyz({2|Z}d~rz$\ !",f&'g-hf(dij2ZkZ2lr[ 0 ,mg3hn4Y==C3==fCG3==C P% Y*SMoN!uwyY{:|:*:*:Y:JY+MNQQU: ! *Y*Y:  W*Y - *Y*Y* Yǵ*,̧:  1#b $)1:BKVu$p !" t$1:K-V$un @ 0 K34py }!*̾!*2**ö*Y:>=52LYY*+:+:*+̱#* '07J\dt|$*}!"J*\d0 \4E1 O*3<=!w*ö**ʶ ! w M*: ,*: ==#2 (/:>EGN$O!" 14} pq*U`*,|N,|+\:-:*<: Bf:GI:JY+MNQU QI:  !: '):  <:  Bf:  -/136 7/136 9;: ?I: AI:!CY SYSYESYESYSYSJ M:!Q _X_-S#f 0:Q["'8EOY_$`Vf`&'`*WDXY00.:&Z[Qij[k\ l] ^f _j .s8f:\@RB8 g0 `Vml` ^m 4 :=l"C3 =l"fCG"=l"fCa=C =l"C OP*cs+Ƕdh:nYp:qYs: tW wW&y:  }W: ! W<: m B:  dh:  :  # 6   W  W _+: *: ) Bn:  :    _Ӳ!*+*: !  8GJ#! !(,38?GL[`n   *48B"W#c$p%$sl]sBs.s8ss: hf_V? L ^ Qf D  q  * cg 0 sl` hnQn 4 (=nq V a=nqCY a=nqC= a=nqC- a=nqC% W*+׸۸ޱ# \]$ i0 i F *#^$  i0  i F *#_$  i0  i "