4/com/hazelcast/mapreduce/impl/task/JobSupervisorjava/lang/ObjectJobSupervisor.javaBcom/hazelcast/mapreduce/impl/task/JobSupervisor$GetResultsRunnableGetResultsRunnable1com/hazelcast/mapreduce/impl/task/JobSupervisor$2 1com/hazelcast/mapreduce/impl/task/JobSupervisor$1 java/util/Map$Entry  java/util/MapEntry/com/hazelcast/mapreduce/JobPartitionState$State)com/hazelcast/mapreduce/JobPartitionStateStateIcom/hazelcast/mapreduce/impl/operation/RequestPartitionResult$ResultState=com/hazelcast/mapreduce/impl/operation/RequestPartitionResult ResultStatereducers$Ljava/util/concurrent/ConcurrentMap;YLjava/util/concurrent/ConcurrentMap;remoteReducerseLjava/util/concurrent/ConcurrentMap;>;context-Ljava/util/concurrent/atomic/AtomicReference;aLjava/util/concurrent/atomic/AtomicReference;keyAssignmentsSLjava/util/concurrent/ConcurrentMap;jobOwnerLcom/hazelcast/nio/Address; ownerNodeZ jobTracker1Lcom/hazelcast/mapreduce/impl/AbstractJobTracker; configuration8Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration;mapReduceService/Lcom/hazelcast/mapreduce/impl/MapReduceService;executorService&Ljava/util/concurrent/ExecutorService;jobProcessInformation=Lcom/hazelcast/mapreduce/impl/task/JobProcessInformationImpl;(Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration;Lcom/hazelcast/mapreduce/impl/AbstractJobTracker;ZLcom/hazelcast/mapreduce/impl/MapReduceService;)V()V 46 7&java/util/concurrent/ConcurrentHashMap9 :7  <  >+java/util/concurrent/atomic/AtomicReference@ A7 !" C $ E *+ G () I ,- K ./ M6com/hazelcast/mapreduce/impl/task/JobTaskConfigurationO getJobOwner()Lcom/hazelcast/nio/Address; QR PS &' UgetName()Ljava/lang/String; WX PY-com/hazelcast/mapreduce/impl/MapReduceService[getExecutorService:(Ljava/lang/String;)Ljava/util/concurrent/ExecutorService; ]^ \_ 01 a*com/hazelcast/mapreduce/impl/MapReduceUtilccreateJobProcessInformation(Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration;Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)Lcom/hazelcast/mapreduce/impl/task/JobProcessInformationImpl; ef dg 23 igetJobId kX Pl-com/hazelcast/mapreduce/impl/task/ReducerTasknX(Ljava/lang/String;Ljava/lang/String;Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)V 4p oq/com/hazelcast/mapreduce/impl/AbstractJobTrackersregisterReducerTask2(Lcom/hazelcast/mapreduce/impl/task/ReducerTask;)V uv twthis1Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;nameLjava/lang/String;jobIdgetMapReduceService1()Lcom/hazelcast/mapreduce/impl/MapReduceService; getJobTracker&()Lcom/hazelcast/mapreduce/JobTracker; startTasks3(Lcom/hazelcast/mapreduce/impl/task/MappingPhase;)V0com/hazelcast/mapreduce/impl/task/MapCombineTask(Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration;Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Lcom/hazelcast/mapreduce/impl/task/MappingPhase;)V 4 registerMapCombineTask5(Lcom/hazelcast/mapreduce/impl/task/MapCombineTask;)V t mappingPhase0Lcom/hazelcast/mapreduce/impl/task/MappingPhase;onNotificationD(Lcom/hazelcast/mapreduce/impl/notification/MapReduceNotification;)VGcom/hazelcast/mapreduce/impl/notification/IntermediateChunkNotification lgetReducerTaskC(Ljava/lang/String;)Lcom/hazelcast/mapreduce/impl/task/ReducerTask; tgetChunk()Ljava/util/Map;  processChunk(Ljava/util/Map;)V o?com/hazelcast/mapreduce/impl/notification/LastChunkNotification lgetPartitionId()I  getSender R .(ILcom/hazelcast/nio/Address;Ljava/util/Map;)V oFcom/hazelcast/mapreduce/impl/notification/ReducingFinishedNotification|(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Lcom/hazelcast/mapreduce/impl/notification/ReducingFinishedNotification;)V 4 $java/util/concurrent/ExecutorServicesubmit3(Ljava/lang/Runnable;)Ljava/util/concurrent/Future; icnILcom/hazelcast/mapreduce/impl/notification/IntermediateChunkNotification; reducerTask/Lcom/hazelcast/mapreduce/impl/task/ReducerTask;lcnALcom/hazelcast/mapreduce/impl/notification/LastChunkNotification;rfnHLcom/hazelcast/mapreduce/impl/notification/ReducingFinishedNotification; notificationALcom/hazelcast/mapreduce/impl/notification/MapReduceNotification;notifyRemoteException3(Lcom/hazelcast/nio/Address;Ljava/lang/Throwable;)V;com/hazelcast/mapreduce/impl/task/JobProcessInformationImplcancelPartitionState 6 collectRemoteAddresses()Ljava/util/Set; cancel8()Lcom/hazelcast/mapreduce/impl/task/TrackableJobFuture; asyncCancelRemoteOperations(Ljava/util/Set;)V java/lang/Thread currentThread()Ljava/lang/Thread;  getStackTrace ()[Ljava/lang/StackTraceElement; java/lang/StringBuilder 7Operation failed on node: append-(Ljava/lang/String;)Ljava/lang/StringBuilder; -(Ljava/lang/Object;)Ljava/lang/StringBuilder; toString X  com/hazelcast/util/ExceptionUtilfixAsyncStackTraceH(Ljava/lang/Throwable;[Ljava/lang/StackTraceElement;Ljava/lang/String;)V 4com/hazelcast/mapreduce/impl/task/TrackableJobFuture setResult(Ljava/lang/Object;)Z  java/util/Set remoteAddress throwableLjava/lang/Throwable; addresses,Ljava/util/Set;Ljava/util/Set;future6Lcom/hazelcast/mapreduce/impl/task/TrackableJobFuture;cancelAndNotify(Ljava/lang/Exception;)Z exceptionLjava/lang/Exception;getConfiguration:()Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration;    unregisterTrackableJobJ(Ljava/lang/String;)Lcom/hazelcast/mapreduce/impl/task/TrackableJobFuture;   tunregisterMapCombineTaskF(Ljava/lang/String;)Lcom/hazelcast/mapreduce/impl/task/MapCombineTask;  t 6 java/lang/StringunregisterReducerTask  t odestroyJobSupervisor4(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)Z  \mapCombineTask2Lcom/hazelcast/mapreduce/impl/task/MapCombineTask; getJobResults7()Ljava/util/Map;get()Ljava/lang/Object; %& A'0com/hazelcast/mapreduce/impl/task/DefaultContext)getReducerFactory*()Lcom/hazelcast/mapreduce/ReducerFactory; +, P-"java/util/concurrent/ConcurrentMap/size 1 02mapSize(I)I 45 d6com/hazelcast/util/MapUtil8createHashMapAdapter(I)Ljava/util/Map; :; 9<entrySet > 0?iterator()Ljava/util/Iterator; AB Cjava/util/IteratorEhasNext()Z GH FInext K& FLgetValue N& Ocom/hazelcast/mapreduce/ReducerQfinalizeReduce S& RTgetKey V& Wput8(Ljava/lang/Object;Ljava/lang/Object;)Ljava/lang/Object; YZ [ requestChunk ] *^finalizeCombiners `6 *areducedResultsLjava/lang/Object;entryJLjava/util/Map$Entry;Ljava/util/Map$Entry;Iresult5Ljava/util/Map;Ljava/util/Map;currentContext2Lcom/hazelcast/mapreduce/impl/task/DefaultContext;getReducerByKey5(Ljava/lang/Object;)Lcom/hazelcast/mapreduce/Reducer;(Ljava/lang/Object;)Lcom/hazelcast/mapreduce/Reducer;&(Ljava/lang/Object;)Ljava/lang/Object; %q 0r&com/hazelcast/mapreduce/ReducerFactoryt newReducer vo uw putIfAbsent yZ 0z beginReduce |6 R} oldReducer!Lcom/hazelcast/mapreduce/Reducer;keyreducergetReducerAddressByKey/(Ljava/lang/Object;)Lcom/hazelcast/nio/Address;com/hazelcast/nio/AddressaddressassignKeyReducerAddress getKeyMember  \ oldAddresscheckAssignedMembersAvailablevalues()Ljava/util/Collection;  0(Ljava/util/Collection;)Z  \0(Ljava/lang/Object;Lcom/hazelcast/nio/Address;)Zequals   oldAssignmentcheckFullyProcessed2(Lcom/hazelcast/mapreduce/JobProcessInformation;)V isOwnerNode H -com/hazelcast/mapreduce/JobProcessInformationgetPartitionStates.()[Lcom/hazelcast/mapreduce/JobPartitionState;  ,[Lcom/hazelcast/mapreduce/JobPartitionState;getState3()Lcom/hazelcast/mapreduce/JobPartitionState$State;   PROCESSED1Lcom/hazelcast/mapreduce/JobPartitionState$State;   getNodeEngine ()Lcom/hazelcast/spi/NodeEngine;  P@com/hazelcast/mapreduce/impl/operation/GetResultOperationFactory'(Ljava/lang/String;Ljava/lang/String;)V 4 com/hazelcast/spi/NodeEngine (Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Lcom/hazelcast/spi/NodeEngine;Lcom/hazelcast/mapreduce/impl/operation/GetResultOperationFactory;Ljava/lang/String;Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Lcom/hazelcast/mapreduce/impl/task/TrackableJobFuture;)V 4 getExecutionService&()Lcom/hazelcast/spi/ExecutionService;  hz:async"com/hazelcast/spi/ExecutionService getExecutorH(Ljava/lang/String;)Lcom/hazelcast/util/executor/ManagedExecutorService;  2com/hazelcast/util/executor/ManagedExecutorService partitionState+Lcom/hazelcast/mapreduce/JobPartitionState;partitionStates nodeEngineLcom/hazelcast/spi/NodeEngine;operationFactoryBLcom/hazelcast/mapreduce/impl/operation/GetResultOperationFactory; jobSupervisorrunnableLjava/lang/Runnable;executionService$Lcom/hazelcast/spi/ExecutionService;executor4Lcom/hazelcast/util/executor/ManagedExecutorService;processInformation/Lcom/hazelcast/mapreduce/JobProcessInformation;getOrCreateContextf(Lcom/hazelcast/mapreduce/impl/task/MapCombineTask;)Lcom/hazelcast/mapreduce/impl/task/DefaultContext;(Lcom/hazelcast/mapreduce/impl/task/MapCombineTask;)Lcom/hazelcast/mapreduce/impl/task/DefaultContext;getCombinerFactory+()Lcom/hazelcast/mapreduce/CombinerFactory;  P^(Lcom/hazelcast/mapreduce/CombinerFactory;Lcom/hazelcast/mapreduce/impl/task/MapCombineTask;)V 4 * compareAndSet'(Ljava/lang/Object;Ljava/lang/Object;)Z  A newContext:Lcom/hazelcast/mapreduce/impl/task/DefaultContext;registerReducerEventInterests(ILjava/util/Set;)V0(ILjava/util/Set;)Vjava/lang/IntegervalueOf(I)Ljava/lang/Integer;  (java/util/concurrent/CopyOnWriteArraySet 7addAll  oldSet partitionIdgetReducerEventInterests(I)Ljava/util/Collection;6(I)Ljava/util/Collection;java/util/CollectiongetJobProcessInformation?()Lcom/hazelcast/mapreduce/impl/task/JobProcessInformationImpl;collectResults((ZLjava/util/Map;Ljava/util/Map$Entry;)VN(ZLjava/util/Map;Ljava/util/Map$Entry;)V rjava/util/Listjava/util/ArrayList  7  Cadd   valuelist$Ljava/util/List;Ljava/util/List; reducedResult mergedResults.()Ljava/util/Set;java/util/HashSet 7 CaddAllFilterJobOwner!(Ljava/util/Set;Ljava/util/Set;)V   getOwner !R " remoteReducerAddresses/(Ljava/util/Set;)V \getGlobalTaskScheduler#()Lcom/hazelcast/spi/TaskScheduler; () *a(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Ljava/util/Set;Lcom/hazelcast/spi/NodeEngine;)V 4, -com/hazelcast/spi/TaskScheduler/execute(Ljava/lang/Runnable;)V 12 03 taskScheduler!Lcom/hazelcast/spi/TaskScheduler;[(Ljava/util/Set;Ljava/util/Set;)VtargetsourceprocessReducerFinished0K(Lcom/hazelcast/mapreduce/impl/notification/ReducingFinishedNotification;)Vjava/lang/Throwable<  getAddress ?R @ checkPartitionReductionCompleted(ILcom/hazelcast/nio/Address;)Z BC D@com/hazelcast/mapreduce/impl/operation/RequestPartitionProcessedFREDUCING H IY(Ljava/lang/String;Ljava/lang/String;ILcom/hazelcast/mapreduce/JobPartitionState$State;)V 4K GLprocessRequestk(Lcom/hazelcast/nio/Address;Lcom/hazelcast/mapreduce/impl/operation/ProcessingOperation;)Ljava/lang/Object; NO \PgetResultStateM()Lcom/hazelcast/mapreduce/impl/operation/RequestPartitionResult$ResultState; RS T SUCCESSFULKLcom/hazelcast/mapreduce/impl/operation/RequestPartitionResult$ResultState; VW Xjava/lang/RuntimeExceptionZ.Could not finalize processing for partitionId \(I)Ljava/lang/StringBuilder; ^ _(Ljava/lang/String;)V 4a [bI(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;Ljava/lang/Throwable;)V d dejava/lang/Errorg sneakyThrow)(Ljava/lang/Throwable;)Ljava/lang/Object; ij k?Lcom/hazelcast/mapreduce/impl/operation/RequestPartitionResult;treducerAddressReducer for partition p not registeredrremove t u 2 tq 0xremoteAddresses access$000 :; |x0x1 access$100b(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)Lcom/hazelcast/mapreduce/impl/MapReduceService; access$200k(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)Lcom/hazelcast/mapreduce/impl/task/JobTaskConfiguration; access$300Y(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;ZLjava/util/Map;Ljava/util/Map$Entry;)V  x2x3 access$400d(Lcom/hazelcast/mapreduce/impl/task/JobSupervisor;)Lcom/hazelcast/mapreduce/impl/AbstractJobTracker; SignatureCodeLineNumberTableLocalVariableTable StackMapTableLocalVariableTypeTable InnerClasses SourceFile!  !"#$%&'()*+,-./0123!45"*8*:Y;=*:Y;?*AYBD*:Y;F*,H*J*+L*N*+TV*+Z`b*+*hj+Z:+m:,oY*rxBYJKL%M0Z5[:\?]E^M_ZbceifoghHyz,-*+()./i{|o}|~/*Nk yz/*Ho yzM*HY*L*+ tuyz n++M*H,N-,M+'+M*H,N-,,,"++M*b Y*,W#*2 xy z{ |*}/~;KUZmH /;ZnyznA*j*N*:*-',۶߻Y+,W @" 9@4AyzA'A 5/  5"*j*M*N*,- -+W    *"yz"   A* mL*H+M*H+N--*H+:*N* W,"o* ",16?4Ayz9}|0'!",#g*D(*L*L.e*=37>=M*=@D:J6M:PRU:,X\WƧ+_M+b,#7*F96 "'M\aqtw|H\cdM$eg"R4h'Pikyz wlm|ik M$ef'Pij|ij$noH*=+sRM,7*L.-*L.+xM*=+,{RN--M,~,BRR"*9=B F *9 HyzHd:pn*F+sM,, yzd'2*F+sM,!*N+M*F+,{N--M,0*.0"**'2yz2d$'H;*N*F& yz~!*F+,{N- -, @ *+*!yz!d!'' *+M,N-66"-2:*LZN*Lm:*L:Y-:*H:*:Y*: :  ¹:   WO65N/01$26371=7E8N9W:c=n>s@tCwDEFGI $Ee{|N\}|WScGn<w3z   yz)*Y*L+M*D,,*D(**LNOQ )yz)!"m A*?sN-&YN*?-{:N-,W8"UVWX0Y5Z8]@^40AyzAhA0 0A0E*?sayzh/*je yzQR/*Vi yzH/*Jm yz  /*Lq yz-y,-X-P\Wa,-X : Y :,-X\W-P  :JM:W- F* uvy+z0{9|H~kux>k d+Myyzy)ykyeg+Myj+YL*?M,J,MN*+-*j M,>6=,2:,#"#*V+#$W+F 9* *03M\m{**%M.yz{*%{%*N'M,+N- Y*+,.4$*%yz%56 %&7,DN-J)-M:*V+$WԱF" '*36*'7yz787978797:;[*LZM*LmN+>6+A:*Eb*N*VGY,-JMQ:UY [YY]`c:*fh lW'mp=mB=>'=EPmprxRE(imrnyz{|w}|qhko'BC\*?sN-%[YYq`sc-,vW-w*?y 7""7?HXZ*\yz\h\o'Kz Kz{:*+}H~z/*NH ~z/*LH ~zP*,-H*~z)kg/*HH ~z2  @@