Skip to content

Commit f0b5381

Browse files
committed
add null checks and some code consolidation
1 parent 17cd4ca commit f0b5381

3 files changed

Lines changed: 76 additions & 31 deletions

File tree

src/main/java/tech/stackable/hadoop/StackableTopologyProvider.java

Lines changed: 61 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -305,12 +305,18 @@ private String resolveListenerEndpoint(GenericKubernetesResource listener) {
305305
Endpoints endpoint = client.endpoints().withName(listenerName).get();
306306
LOG.debug("Matched ingressAddress [{}]", listenerName);
307307

308-
if (endpoint.getSubsets().isEmpty()) {
308+
List<EndpointSubset> subsets = endpoint != null ? endpoint.getSubsets() : null;
309+
if (subsets == null || subsets.isEmpty()) {
309310
LOG.warn("Endpoint {} has no subsets - pod may be restarting", listenerName);
310311
return listenerName;
311312
}
312313

313-
EndpointAddress address = endpoint.getSubsets().get(0).getAddresses().get(0);
314+
List<EndpointAddress> addresses = subsets.get(0).getAddresses();
315+
if (addresses == null || addresses.isEmpty()) {
316+
LOG.warn("Endpoint {} has no ready addresses - pod may not be ready", listenerName);
317+
return listenerName;
318+
}
319+
EndpointAddress address = addresses.get(0);
314320
LOG.info(
315321
"Resolved listener {} to IP {} on node {}",
316322
listenerName,
@@ -366,9 +372,12 @@ private void cachePodByNameAndIps(Pod pod) {
366372
LOG.debug("Refreshing pod cache: adding {}", podName);
367373
cache.putPod(podName, pod);
368374

369-
// Cache by all IPs - this is crucial for IP-based lookups
370-
for (PodIP ip : pod.getStatus().getPodIPs()) {
371-
cache.putPod(ip.getIp(), pod);
375+
PodStatus podStatus = pod.getStatus();
376+
if (podStatus != null && podStatus.getPodIPs() != null) {
377+
// Cache by all IPs - this is crucial for IP-based lookups
378+
for (PodIP ip : podStatus.getPodIPs()) {
379+
cache.putPod(ip.getIp(), pod);
380+
}
372381
}
373382
}
374383

@@ -421,7 +430,7 @@ private String resolveToIpAddress(String hostname) {
421430
LOG.debug("Resolved {} to {}", hostname, ip);
422431
return ip;
423432
} catch (UnknownHostException e) {
424-
LOG.warn("Failed to resolve address: {} - defaulting to {}", hostname, DEFAULT_RACK);
433+
LOG.warn("Failed to resolve address: {}", hostname);
425434
return hostname;
426435
}
427436
}
@@ -439,12 +448,16 @@ private Map<String, String> buildNodeToDatanodeMap(List<Pod> dataNodes) {
439448

440449
for (Pod dataNode : dataNodes) {
441450
String nodeName = dataNode.getSpec().getNodeName();
442-
String dataNodeIp = dataNode.getStatus().getPodIP();
451+
PodStatus podStatus = dataNode.getStatus();
452+
String dataNodeIp = podStatus != null ? podStatus.getPodIP() : null;
443453

444454
if (nodeName != null && dataNodeIp != null) {
455+
Node node = getOrFetchNode(nodeName);
456+
if (node == null) {
457+
continue;
458+
}
445459
LOG.debug("Assigned to node-name [{}/{}]", nodeName, dataNodeIp);
446460
nodeToDatanode.put(nodeName, dataNodeIp);
447-
Node node = getOrFetchNode(nodeName);
448461
for (NodeAddress nodeAddress : node.getStatus().getAddresses()) {
449462
LOG.debug("Assigned to node-address [{}/{}]", nodeAddress.getAddress(), dataNodeIp);
450463
nodeToDatanode.put(nodeAddress.getAddress(), dataNodeIp);
@@ -513,11 +526,14 @@ private Map<String, Map<String, String>> buildNodeLabelMap(List<Pod> dataNodes)
513526
String nodeName = dataNode.getSpec().getNodeName();
514527

515528
if (nodeName == null) {
516-
LOG.warn("Pod [{}] not yet assigned to node, retrying", dataNode.getMetadata().getName());
517-
return result;
529+
LOG.warn("Pod [{}] not yet assigned to node...", dataNode.getMetadata().getName());
530+
continue;
518531
}
519532

520533
Node node = getOrFetchNode(nodeName);
534+
if (node == null) {
535+
continue;
536+
}
521537
Map<String, String> nodeLabels = node.getMetadata().getLabels();
522538
LOG.debug("Labels for node [{}]:[{}]....", nodeName, nodeLabels);
523539

@@ -526,9 +542,12 @@ private Map<String, Map<String, String>> buildNodeLabelMap(List<Pod> dataNodes)
526542
result.put(nodeAddress.getAddress(), nodeLabels);
527543
}
528544

529-
for (PodIP podIp : dataNode.getStatus().getPodIPs()) {
530-
LOG.debug("...assigned to IP [{}]", podIp.getIp());
531-
result.put(podIp.getIp(), nodeLabels);
545+
PodStatus podStatus = dataNode.getStatus();
546+
if (podStatus != null && podStatus.getPodIPs() != null) {
547+
for (PodIP podIp : podStatus.getPodIPs()) {
548+
LOG.debug("...assigned to IP [{}]", podIp.getIp());
549+
result.put(podIp.getIp(), nodeLabels);
550+
}
532551
}
533552
}
534553
return result;
@@ -539,6 +558,10 @@ private Node getOrFetchNode(String nodeName) {
539558
if (node == null) {
540559
LOG.debug("Fetching node: {}", nodeName);
541560
node = client.nodes().withName(nodeName).get();
561+
if (node == null) {
562+
LOG.warn("Node {} not found in cluster", nodeName);
563+
return null;
564+
}
542565
cache.putNode(nodeName, node);
543566
}
544567
return node;
@@ -558,7 +581,11 @@ private Map<String, Map<String, String>> buildPodLabelMap(List<Pod> dataNodes) {
558581
Map<String, String> podLabels = pod.getMetadata().getLabels();
559582
LOG.debug("Labels for pod [{}]:[{}]....", pod.getMetadata().getName(), podLabels);
560583

561-
for (PodIP podIp : pod.getStatus().getPodIPs()) {
584+
PodStatus podStatus = pod.getStatus();
585+
if (podStatus == null || podStatus.getPodIPs() == null) {
586+
continue;
587+
}
588+
for (PodIP podIp : podStatus.getPodIPs()) {
562589
LOG.debug("...assigned to pod IP [{}]", podIp.getIp());
563590
result.put(podIp.getIp(), podLabels);
564591
}
@@ -578,27 +605,34 @@ private void startPodInformer() {
578605
new ResourceEventHandler<>() {
579606
@Override
580607
public void onAdd(Pod pod) {
581-
cache.putPod(pod.getMetadata().getName(), pod);
582-
for (PodIP ip : pod.getStatus().getPodIPs()) {
583-
cache.putPod(ip.getIp(), pod);
584-
}
585-
LOG.debug("Pod {} added", pod.getMetadata().getName());
608+
addPod(pod);
609+
}
610+
611+
@Override
612+
public void onDelete(Pod pod, boolean deletedFinalStateUnknown) {
613+
deletePod(pod);
586614
}
587615

588616
@Override
589617
public void onUpdate(Pod oldPod, Pod newPod) {
590-
cache.putPod(oldPod.getMetadata().getName(), newPod);
591-
for (PodIP ip : oldPod.getStatus().getPodIPs()) {
592-
cache.putPod(ip.getIp(), newPod);
593-
}
618+
// In case the IPs are updated, update all IP keys
619+
deletePod(oldPod);
620+
addPod(newPod);
594621
LOG.trace("Pod {} updated", oldPod.getMetadata().getName());
595622
}
596623

597-
@Override
598-
public void onDelete(Pod pod, boolean deletedFinalStateUnknown) {
624+
private void addPod(Pod pod) {
625+
cachePodByNameAndIps(pod);
626+
LOG.debug("Pod {} added", pod.getMetadata().getName());
627+
}
628+
629+
private void deletePod(Pod pod) {
599630
cache.deletePod(pod.getMetadata().getName());
600-
for (PodIP ip : pod.getStatus().getPodIPs()) {
601-
cache.deletePod(ip.getIp());
631+
PodStatus podStatus = pod.getStatus();
632+
if (podStatus != null && podStatus.getPodIPs() != null) {
633+
for (PodIP ip : podStatus.getPodIPs()) {
634+
cache.deletePod(ip.getIp());
635+
}
602636
}
603637
LOG.debug("Pod {} deleted", pod.getMetadata().getName());
604638
}

src/main/java/tech/stackable/hadoop/TopologyUtils.java

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package tech.stackable.hadoop;
22

33
import io.fabric8.kubernetes.api.model.GenericKubernetesResource;
4+
import java.util.Collections;
45
import java.util.List;
56
import java.util.Map;
67
import java.util.stream.Collectors;
@@ -16,11 +17,20 @@ public class TopologyUtils {
1617

1718
public static List<String> getIngressAddresses(GenericKubernetesResource listener) {
1819
// suppress warning as we know the structure of our own listener resource
20+
Object statusObj = listener.getAdditionalProperties().get(STATUS);
21+
if (statusObj == null) {
22+
LOG.warn("Listener {} has no status", listener.getMetadata().getName());
23+
return Collections.emptyList();
24+
}
25+
@SuppressWarnings("unchecked")
26+
Map<String, Object> status = (Map<String, Object>) statusObj;
27+
Object addressesObj = status.get(INGRESS_ADDRESSES);
28+
if (addressesObj == null) {
29+
LOG.warn("Listener {} has no ingress addresses", listener.getMetadata().getName());
30+
return Collections.emptyList();
31+
}
1932
@SuppressWarnings("unchecked")
20-
List<Map<String, Object>> ingressAddresses =
21-
((List<Map<String, Object>>)
22-
((Map<String, Object>) listener.getAdditionalProperties().get(STATUS))
23-
.get(INGRESS_ADDRESSES));
33+
List<Map<String, Object>> ingressAddresses = (List<Map<String, Object>>) addressesObj;
2434
return ingressAddresses.stream()
2535
.map(ingress -> (String) ingress.get(ADDRESS))
2636
.collect(Collectors.toList());

test/topology-provider/stack/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,3 +15,4 @@ kubectl apply -f ./hdfs-utils/test/topology-provider/stack/03-hdfs.yaml
1515
kubectl apply -f ./hdfs-utils/test/topology-provider/stack/04-spark.yaml
1616
kubectl apply -f ./hdfs-utils/test/topology-provider/stack/05-access-hdfs.yaml
1717
```
18+

0 commit comments

Comments
 (0)