@@ -119,24 +119,30 @@ private List<String> tryResolveFromCache(List<String> names) {
119119 private List <String > performFullResolution (List <String > names ) {
120120 LOG .debug ("Performing full topology resolution for: {}" , names );
121121
122- // Step 1: Gather all dataNodes
122+ // Pre-requisites : fetch all dataNodes and build label lookup maps from them
123123 List <Pod > dataNodes = fetchDataNodes ();
124-
125- // Step 2: Resolve listeners to actual datanode IPs
126- List <String > resolvedNames = resolveListeners (names );
127-
128- // Step 3: Build label lookup maps
129124 Map <String , Map <String , String >> podLabels = buildPodLabelMap (dataNodes );
130125 Map <String , Map <String , String >> nodeLabels = buildNodeLabelMap (dataNodes );
131126
132- // Step 4: Build node-to-datanode map for O(1) colocated lookups
127+ // Build node-to-datanode map for O(1) colocated lookups
133128 Map <String , String > nodeToDatanodeIp = buildNodeToDatanodeMap (dataNodes );
134129
135- // Step 5: Resolve client pods to co-located dataNodes
130+ // Resolve masqueraded IPs to nodes: do this before inspecting possible listener-
131+ // or other pod-IPs as we don't want to mistakenly treat a masqueraded IP as a
132+ // cache-miss
133+ List <String > resolvedNames = tryResolveNodes (names , nodeToDatanodeIp );
134+
135+ // Resolve dataNode listeners to datanode IPs
136+ resolvedNames = tryResolveListeners (resolvedNames );
137+
138+ // Step : Resolve client pods to co-located dataNodes
136139 List <String > datanodeIps =
137- resolveClientPodsToDataNodes (resolvedNames , podLabels , nodeToDatanodeIp );
140+ tryResolveClientPodsToDataNodes (resolvedNames , podLabels , nodeToDatanodeIp );
138141
139- // Step 6: Build topology strings and cache results
142+ // Step : Build topology strings and cache results
143+ // IP-masquerading can mean that the advertised IP is either a client Pod's node IP,
144+ // or an IP that is nto easily associated with a node (e.g. if the veth interface is
145+ // used).
140146 return buildAndCacheTopology (names , datanodeIps , podLabels , nodeLabels );
141147 }
142148
@@ -183,6 +189,25 @@ private List<Pod> fetchDataNodes() {
183189 return dataNodes ;
184190 }
185191
192+ // ============================================================================
193+ // NODE RESOLUTION
194+ // ============================================================================
195+
196+ private List <String > tryResolveNodes (List <String > names , Map <String , String > nodeToDatanodeIp ) {
197+ List <String > result = new ArrayList <>();
198+
199+ for (String name : names ) {
200+ String dataNodeIp = nodeToDatanodeIp .get (name );
201+ if (dataNodeIp == null ) {
202+ result .add (name );
203+ } else {
204+ LOG .debug ("Returning dataNode {} for {}" , name , dataNodeIp );
205+ result .add (dataNodeIp );
206+ }
207+ }
208+ return result ;
209+ }
210+
186211 // ============================================================================
187212 // LISTENER RESOLUTION
188213 // ============================================================================
@@ -221,10 +246,10 @@ private String getListenerVersion() {
221246 }
222247 }
223248
224- private List <String > resolveListeners (List <String > names ) {
249+ private List <String > tryResolveListeners (List <String > names ) {
225250 refreshListenerCacheIfNeeded (names );
226251
227- return names .stream ().map (this ::resolveListenerToDatanode ).collect (Collectors .toList ());
252+ return names .stream ().map (this ::tryResolveListenerToDatanode ).collect (Collectors .toList ());
228253 }
229254
230255 private void refreshListenerCacheIfNeeded (List <String > names ) {
@@ -268,7 +293,7 @@ private void cacheListenerByNameAndAddresses(GenericKubernetesResource listener)
268293 * listener
269294 * @return either the name (for non-listener) or the dataNode IP to which this listener resolves
270295 */
271- private String resolveListenerToDatanode (String name ) {
296+ private String tryResolveListenerToDatanode (String name ) {
272297 GenericKubernetesResource listener = cache .getListener (name );
273298 if (listener == null ) {
274299 LOG .debug ("Not a listener: {}" , name );
@@ -314,7 +339,7 @@ private GenericKubernetesResourceList fetchListeners(String listenerVersion) {
314339 // CLIENT POD RESOLUTION
315340 // ============================================================================
316341
317- private List <String > resolveClientPodsToDataNodes (
342+ private List <String > tryResolveClientPodsToDataNodes (
318343 List <String > names ,
319344 Map <String , Map <String , String >> podLabels ,
320345 Map <String , String > nodeToDatanodeIp ) {
@@ -408,7 +433,8 @@ private String resolveToIpAddress(String hostname) {
408433
409434 /**
410435 * Build a map from Kubernetes node name to datanode IP. This enables O(1) lookup when finding
411- * co-located dataNodes for client pods.
436+ * co-located dataNodes for client pods. This map will contain as keys both the node name and all
437+ * its addresses (as the address may be used by pods with IP masquerading).
412438 *
413439 * <p>Note: If multiple dataNodes run on the same node, the last one wins. This is acceptable
414440 * because all dataNodes on the same node have the same topology.
@@ -421,11 +447,17 @@ private Map<String, String> buildNodeToDatanodeMap(List<Pod> dataNodes) {
421447 String dataNodeIp = dataNode .getStatus ().getPodIP ();
422448
423449 if (nodeName != null && dataNodeIp != null ) {
450+ LOG .debug ("Assigned to node-name [{}/{}]" , nodeName , dataNodeIp );
424451 nodeToDatanode .put (nodeName , dataNodeIp );
452+ Node node = getOrFetchNode (nodeName );
453+ for (NodeAddress nodeAddress : node .getStatus ().getAddresses ()) {
454+ LOG .debug ("Assigned to node-address [{}/{}]" , nodeAddress .getAddress (), dataNodeIp );
455+ nodeToDatanode .put (nodeAddress .getAddress (), dataNodeIp );
456+ }
425457 }
426458 }
427459
428- LOG .debug ("Built node-to-datanode map with {} entries " , nodeToDatanode . size () );
460+ LOG .debug ("Built node-to-datanode map {} " , nodeToDatanode );
429461 return nodeToDatanode ;
430462 }
431463
@@ -474,27 +506,32 @@ private String extractLabelValue(
474506 * Given a list of dataNodes this function will resolve which dataNodes run on which node as well
475507 * as all the ips assigned to a dataNodes. It will then return a mapping of every ip address to
476508 * the labels that are attached to the "physical" node running the dataNodes that this ip belongs
477- * to.
509+ * to. It will also do this for the node addresses as calling pods may masquerade as node IPs.
478510 *
479511 * @param dataNodes List of all in-scope dataNodes (datanode pods in this namespace)
480512 * @return Map of ip addresses to labels of the node running the pod that the ip address belongs
481513 * to
482514 */
483515 private Map <String , Map <String , String >> buildNodeLabelMap (List <Pod > dataNodes ) {
484516 Map <String , Map <String , String >> result = new HashMap <>();
485- for (Pod pod : dataNodes ) {
486- String nodeName = pod .getSpec ().getNodeName ();
517+ for (Pod dataNode : dataNodes ) {
518+ String nodeName = dataNode .getSpec ().getNodeName ();
487519
488520 if (nodeName == null ) {
489- LOG .warn ("Pod [{}] not yet assigned to node, retrying" , pod .getMetadata ().getName ());
521+ LOG .warn ("Pod [{}] not yet assigned to node, retrying" , dataNode .getMetadata ().getName ());
490522 return result ;
491523 }
492524
493525 Node node = getOrFetchNode (nodeName );
494526 Map <String , String > nodeLabels = node .getMetadata ().getLabels ();
495527 LOG .debug ("Labels for node [{}]:[{}]...." , nodeName , nodeLabels );
496528
497- for (PodIP podIp : pod .getStatus ().getPodIPs ()) {
529+ for (NodeAddress nodeAddress : node .getStatus ().getAddresses ()) {
530+ LOG .debug ("...assigned to node address [{}]" , nodeAddress .getAddress ());
531+ result .put (nodeAddress .getAddress (), nodeLabels );
532+ }
533+
534+ for (PodIP podIp : dataNode .getStatus ().getPodIPs ()) {
498535 LOG .debug ("...assigned to IP [{}]" , podIp .getIp ());
499536 result .put (podIp .getIp (), nodeLabels );
500537 }
0 commit comments