Github user mxm commented on a diff in the pull request:

    https://github.com/apache/flink/pull/2657#discussion_r84500666
  
    --- Diff: 
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java
 ---
    @@ -202,101 +205,125 @@ public void shutDown() throws Exception {
        //  RPC methods
        // 
------------------------------------------------------------------------
     
    -   /**
    -    * Register a {@link JobMaster} at the resource manager.
    -    *
    -    * @param resourceManagerLeaderId The fencing token for the 
ResourceManager leader
    -    * @param jobMasterAddress        The address of the JobMaster that 
registers
    -    * @param jobID                   The Job ID of the JobMaster that 
registers
    -    * @return Future registration response
    -    */
        @RpcMethod
    -   public Future<RegistrationResponse> registerJobMaster(
    -           final UUID resourceManagerLeaderId, final UUID 
jobMasterLeaderId,
    -           final String jobMasterAddress, final JobID jobID) {
    +   public Future<RegistrationResponse> registerJobManager(
    +                   final UUID resourceManagerLeaderId,
    +                   final UUID jobManagerLeaderId,
    +                   final String jobManagerAddress,
    +                   final JobID jobId) {
    +
    +           checkNotNull(resourceManagerLeaderId);
    +           checkNotNull(jobManagerLeaderId);
    +           checkNotNull(jobManagerAddress);
    +           checkNotNull(jobId);
    +
    +           if (isValid(resourceManagerLeaderId)) {
    +                   if (!jobLeaderIdService.containsJob(jobId)) {
    +                           try {
    +                                   jobLeaderIdService.addJob(jobId);
    +                           } catch (Exception e) {
    +                                   // This should actually never happen 
because, it should always be possible to add a new job
    +                                   ResourceManagerException exception = 
new ResourceManagerException("Could not add the job " +
    +                                           jobId + " to the job id leader 
service. This should never happen.", e);
    +
    +                                   onFatalErrorAsync(exception);
    +
    +                                   log.debug("Could not add job {} to job 
leader id service.", jobId, e);
    +                                   return 
FlinkCompletableFuture.completedExceptionally(exception);
    +                           }
    +                   }
     
    -           checkNotNull(jobMasterAddress);
    -           checkNotNull(jobID);
    +                   log.info("Registering job manager {}@{} for job {}.", 
jobManagerLeaderId, jobManagerAddress, jobId);
    +
    +                   Future<UUID> jobLeaderIdFuture;
     
    -           // create a leader retriever in case it doesn't exist
    -           final JobIdLeaderListener jobIdLeaderListener;
    -           if (leaderListeners.containsKey(jobID)) {
    -                   jobIdLeaderListener = leaderListeners.get(jobID);
    -           } else {
                        try {
    -                           LeaderRetrievalService jobMasterLeaderRetriever 
=
    -                                   
highAvailabilityServices.getJobManagerLeaderRetriever(jobID);
    -                           jobIdLeaderListener = new 
JobIdLeaderListener(jobID, jobMasterLeaderRetriever);
    +                           jobLeaderIdFuture = 
jobLeaderIdService.getLeaderId(jobId);
                        } catch (Exception e) {
    -                           log.warn("Failed to start 
JobMasterLeaderRetriever for job id {}", jobID, e);
    +                           // we cannot check the job leader id so let's 
fail
    +                           // TODO: Maybe it's also ok to skip this check 
in case that we cannot check the leader id
    +                           ResourceManagerException exception = new 
ResourceManagerException("Cannot obtain the " +
    +                                   "job leader id future to verify the 
correct job leader.", e);
    +
    +                           onFatalErrorAsync(exception);
    --- End diff --
    
    Declining seems ok in this case since the failure might be temporary.


---
If your project is set up for it, you can reply to this email and have your
reply appear on GitHub as well. If your project does not have this feature
enabled and wishes so, or if the feature is enabled but not working, please
contact infrastructure at [email protected] or file a JIRA ticket
with INFRA.
---

Reply via email to