【问题标题】:How to track stateful pods created in a k8s cluster?如何跟踪在 k8s 集群中创建的有状态 Pod?
【发布时间】:2020-12-14 23:15:35
【问题描述】:

设置

我通过以下方式设置了 k8s 集群:

  • 1 个主节点
  • 2 个工作节点

集群是使用 kubeadm 和 Flannel 设置的。

我有两种不同的 pod 类型:

  • Java 代理服务器
  • Java TCP 服务器

Java Proxy Server pod 是我最初创建的 StatefulSet。每个 Java 代理服务器都有自己的状态(当前连接的客户端),但是它们都应该共享一个公共状态。

此常见状态是 Java TCP Server pod 及其关联 IP 地址的最新列表。 我的目标是确保每个代理服务器都有一个当前的 TCP 服务器列表,它可以代理连接。

Java TCP 服务器的每个实例都有自己独特的状态,并且也部署为 StatefulSet。 TCP 服务器 pod 之间唯一的共同点是它们可以接收来自代理服务器的连接。

代理服务器必须知道 TCP 服务器 pod 何时启动或关闭,以便他们知道哪些 pod 可用于代理连接。

TCP 服务器由代理服务器委托连接。绝不会出现代理服务器随机给 TCP 服务器一个连接并且它们没有负载平衡的情况。

尝试

我尝试使用Java Kubernetes Client 并在我的代理服务器上实现了一个监视,如下所示:

ApiClient apiClient = Config.defaultClient();
apiClient.setReadTimeout(0);
System.out.println(apiClient.getBasePath());
Configuration.setDefaultApiClient(apiClient);

CoreV1Api api = new CoreV1Api();
V1PodList pods = api.listPodForAllNamespaces(null, null, null, null, null, null, null, null, null);
V1ListMeta podsMeta = pods.getMetadata();
if (podsMeta != null) {
    String resourceVersion = podsMeta.getResourceVersion();

    Watch<V1Pod> watch = Watch.createWatch(
            apiClient,
            api.listPodForAllNamespacesCall(null, null, null, null, null, null, resourceVersion, null, true, null),
            new TypeToken<Watch.Response<V1Pod>>(){}.getType());

    while (watch.hasNext()) {
        Watch.Response<V1Pod> response = watch.next();
        V1Pod pod = response.object;
        V1PodStatus status = pod.getStatus();
        if (status != null) {
            System.out.printf("Pod IP: %s\n", status.getPodIP());
            System.out.printf("Pod Reason: %s\n", status.getReason());
        }
    }

    watch.close();
}

这工作相对较好。对我来说最大的问题是,对于这个简单的过程,它为我的最终 Jar 文件增加了 40MB 的巨大空间。

我知道 40MB 对某些人来说可能不算多。我只是觉得有一种更轻量级的方式来实现我正在尝试做的事情?

是否有更好的流程来跟踪在我忽略的集群中创建和销毁的这些 pod?

【问题讨论】:

  • 不知道这是否有用,但是this article关于如何通过statefulset配置自愈redis集群可能会提供一些建议。
  • 这听起来像是 k8s 负载均衡器服务将为您做的事情。您需要自己的代理有什么原因吗? Kubernetes 将使用 pod liveness/readiness 探测来了解怎么了可用的。
  • @aguest 是的,不幸的是我的项目设置方式我需要有自己的代理。除了代理来自我的 TCP 服务器的数据包之外,“代理”本身会拦截客户端数据包并对其进行修改以及发送自己的数据包。

标签: java kubernetes proxy kubernetes-statefulset


【解决方案1】:

我已经提出了自己的解决方案(目前),虽然不是很漂亮。我决定使用边车模式,并在我的服务器代理 pod 中附带另一个容器。它是用 Go 编写的,二进制文件被剥离并在 Alpine Linux 上运行。

目前,我只是使用简单的 UDP 连接到 Java 代理服务器,让它知道何时删除或添加了 pod。

我引用了Kubernetes docs,它描述了如何使用 kube-apiserver 服务。就我而言,我只是使用默认服务帐户。

我将附上一些示例代码来说明实现。基本上所有 API 凭据都是由 Kubernetes 默认通过文件提供的。

/var/run/secrets/kubernetes.io/serviceaccount/token /var/run/secrets/kubernetes.io/serviceaccount/ca.crt

然后我们可以使用这些凭据通过默认 DNS 配置访问 Kubernetes 提供的 API。

https://kubernetes.default.svc/api/v1/

以下是关于我如何能够在集群中跟踪 pod:

package main

import (
    "context"
    "crypto/tls"
    "crypto/x509"
    "encoding/json"
    "fmt"
    "io/ioutil"
    "log"
    "net/http"
)

func main() {
    backGroundContext := context.Background()

    accessTokenData, accessTokenFileError := ioutil.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/token")
    if accessTokenFileError != nil {
        log.Fatalln(accessTokenFileError)
    }

    accessToken := string(accessTokenData)

    k8sCertificate, certificateFileError := ioutil.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/ca.crt")
    if certificateFileError != nil {
        log.Fatalln(certificateFileError)
    }

    certificateAuthorityPool := x509.NewCertPool()
    certificateAuthorityPool.AppendCertsFromPEM(k8sCertificate)

    client := &http.Client{
        Transport: &http.Transport{
            TLSClientConfig: &tls.Config{
                RootCAs: certificateAuthorityPool,
            },
        },
        Timeout: 0, // disable timeout for the watch request
    }

    request, requestError := http.NewRequestWithContext(backGroundContext, "GET", "https://kubernetes.default.svc/api/v1/namespaces/default/pods", nil)
    if requestError != nil {
        log.Fatalln(requestError)
    }
    request.Header.Add("Authorization", fmt.Sprintf("Bearer %s", accessToken))

    response, responseError := client.Do(request)
    if responseError != nil {
        log.Fatalln(responseError)
    }
    defer response.Body.Close()

    type PodListMetaData struct {
        ResourceVersion string `json:"resourceVersion"`
    }

    type PodStatus struct {
        IP string `json:"podIP"`
    }

    type PodMetaData struct {
        Name string `json:"name"`
    }

    type PodResult struct {
        MetaData PodMetaData `json:"metadata"`
        Status   PodStatus   `json:"status"`
    }

    type PodListResult struct {
        MetaData PodListMetaData `json:"metadata"`
        Items    []PodResult     `json:"items"`
    }

    var list PodListResult

    decoder := json.NewDecoder(response.Body)
    decodeError := decoder.Decode(&list)
    if decodeError != nil {
        log.Fatalln(decodeError)
    }

    resourceVersion := list.MetaData.ResourceVersion
    log.Printf("Resource Version: %s\n", resourceVersion)

    for _, item := range list.Items {
        log.Printf("Found Pod: %s with IP of %s\n", item.MetaData.Name, item.Status.IP)
    }

    watchRequest, watchRequestError := http.NewRequestWithContext(backGroundContext, "GET", fmt.Sprintf("https://kubernetes.default.svc/api/v1/namespaces/default/pods?watch=1&resourceVersion=%s&allowWatchBookmarks=true", resourceVersion), nil)
    if watchRequestError != nil {
        log.Fatalln(watchRequestError)
    }
    watchRequest.Header.Add("Authorization", fmt.Sprintf("Bearer %s", accessToken))

    response1, response1Error := client.Do(watchRequest)
    if response1Error != nil {
        log.Fatalln(response1Error)
    }
    defer response1.Body.Close()

    type PodListWatchResult struct {
        Type   string    `json:"type"`
        Object PodResult `json:"object"`
    }

    decoder1 := json.NewDecoder(response1.Body)

    for decoder1.More() {
        var podResult PodListWatchResult
        decodeError1 := decoder1.Decode(&podResult)
        if decodeError1 != nil {
            log.Fatalln(decodeError1)
        }

        log.Printf("Found Pod: %s with IP of %s\n", podResult.Object.MetaData.Name, podResult.Object.Status.IP)
    }
}

这里完全有改进的余地,但这是我的解决方案的要点。每次流式传输新的 JSON 对象时,我都会对其进行处理以确保它是我感兴趣的内容,然后通过 UDP 将其转发到 Java 代理服务器。

基于 Alpine 构建的二进制文件大约 10 MB,非常轻!

我怀疑“真正”解决方案的唯一选择是使用官方的Kubernetes Go ClientJava Client。它会使我的容器更大,但我的解决方案没有太多保证,充其量看起来像是一个 hack。

我仍然希望有一些我忽略的东西可以简化这一切并且不需要大型客户端库。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-07-10
    • 1970-01-01
    • 2021-11-27
    • 2019-09-25
    • 2018-10-29
    • 2019-06-28
    • 2019-12-20
    相关资源
    最近更新 更多