Ansible

介绍

Ansible是一种开源的自动化运维工具,可以实现配置管理、应用部署、任务执行等操作,通过SSH协议进行通信,无需在远程主机上安装客户端,具有简单易用、扩展性强、可靠稳定等特点。

pip install ansible

Inventory

Inventory是一个用于描述被管理节点的列表文件。它定义了Ansible任务将要在哪些主机上执行。清单可以包含IP地址、主机名、组名等信息,以及变量和组相关的配置。

清单文件通常采用INI格式或YAML格式编写,其中每个主机都有对应的主机名或IP地址,并且可以根据需求进行分组。通过在清单中定义主机组,可以更方便地对不同组的主机进行管理和操作。

除了静态清单外,Ansible还支持动态清单,这意味着清单可以从外部数据源(如云平台、自定义脚本、数据库等)动态生成。这样可以实现按需扩展或动态调整被管理节点,提高了灵活性和可扩展性。

使用清单,您可以轻松地指定要在哪些主机上执行特定的Ansible任务或剧本,使得管理和组织被管理节点变得简单而高效。

[self]
127.0.0.1 var1=value1

[all:vars]
all_var=all

Ad-HOC

Ansible的Ad-Hoc命令是一种快速、临时性的方式来执行单个任务或命令,而无需编写和执行完整的Playbook。通过Ad-Hoc命令,可以直接在终端上与远程主机进行交互和操作。

ansible [pattern] -m [module] -a "[module options]"

常用参数

位置参数:
  pattern       匹配主机模式

选项:
  --list-hosts          输出匹配的主机列表,不执行其他任何操作
  --playbook-dir BASEDIR
                        由于此工具不使用playbooks,请将其用作替代playbook目录。这设置了许多功能的相对路径,包括roles/ group_vars/等。
  --task-timeout TASK_TIMEOUT
                        设置任务的超时限制(以秒为单位),必须是正整数。
  -B SECONDS, --background SECONDS
                        异步运行,在X秒后失败(默认=N/A)
  -M MODULE_PATH, --module-path MODULE_PATH
                        在模块库之前添加以冒号分隔的路径(默认={{ ANSIBLE_HOME ~ "/plugins/modules:/usr/share/ansible/plugins/modules" }})
  -P POLL_INTERVAL, --poll POLL_INTERVAL
                        如果使用-B,设置轮询间隔(默认=15)
  -a MODULE_ARGS, --args MODULE_ARGS
                        使用空格分隔的k=v格式指定操作的选项:-a 'opt1=val1 opt2=val2' 或使用JSON字符串:-a '{"opt1": "val1", "opt2": "val2"}'
  -e EXTRA_VARS, --extra-vars EXTRA_VARS
                        设置额外的变量,格式为key=value或YAML/JSON,如果是文件名,请在前面加上@
  -f FORKS, --forks FORKS
                        指定要使用的并行进程数(默认=5)
  -h, --help            显示此帮助消息并退出
  -i INVENTORY, --inventory INVENTORY, --inventory-file INVENTORY
                        指定清单主机路径或逗号分隔的主机列表。--inventory-file已废弃
  -l SUBSET, --limit SUBSET
                        进一步限制所选主机到附加的模式
  -m MODULE_NAME, --module-name MODULE_NAME
                        要执行的操作的名称(默认=command)
  -o, --one-line        精简输出
  -v, --verbose         导致Ansible打印更多的调试消息。添加多个-v将增加详细程度,内置插件当前可评估至-vvvvvv。开始时一个合理的级别是-vvv,连接调试可能需要-vvvv。

特权升级选项:
  --become-method BECOME_METHOD
                        要使用的特权升级方法(默认=sudo),使用ansible-doc -t become -l列出有效选项。
  --become-user BECOME_USER
                        以此用户身份运行操作(默认=root)
  -b, --become          使用特权升级运行操作(不会提示密码)

连接选项:
  --private-key PRIVATE_KEY_FILE, --key-file PRIVATE_KEY_FILE
                        使用此文件进行连接身份验证
  --scp-extra-args SCP_EXTRA_ARGS
                        指定传递给scp的额外参数(例如-l)
  --sftp-extra-args SFTP_EXTRA_ARGS
                        指定传递给sftp的额外参数(例如-f、-l)
  --ssh-common-args SSH_COMMON_ARGS
                        指定传递给sftp/scp/ssh的常见参数(例如ProxyCommand)
  --ssh-extra-args SSH_EXTRA_ARGS
                        指定传递给ssh的额外参数(例如-R)
  -T TIMEOUT, --timeout TIMEOUT
                        覆盖连接超时时间(以秒为单位)(默认=10)
  -c CONNECTION, --connection CONNECTION
                        要使用的连接类型(默认=smart)
  -u REMOTE_USER, --user REMOTE_USER
                        以此用户身份连接(默认=None)

Playbook

Ansible的Playbook是一种用于定义、配置和部署多个任务的文件。它是一种结构化的、可扩展的自动化工具,可以实现复杂的运维流程和长期管理。

Playbook使用YAML格式编写,由一个或多个任务(tasks)组成。每个任务都定义了要在目标主机上执行的具体操作,如执行命令、复制文件、安装软件包等。任务可以根据需要指定不同的模块、参数和条件。

ansible-playbook playbook [playbook ...]

以下是一个简单的Playbook示例:

- name: Install and start Nginx
  hosts: web_servers
  become: true

  tasks:
    - name: Install Nginx package
      apt:
        name: nginx
        state: present

    - name: Copy Nginx configuration file
      copy:
        src: /path/to/nginx.conf
        dest: /etc/nginx/nginx.conf
        owner: root
        group: root
        mode: '0644'

    - name: Start Nginx service
      service:
        name: nginx
        state: started

这个Playbook中,name字段是Playbook的名称,hosts字段指定了目标主机或主机组,become: true表示以特权用户身份运行任务。

通过执行ansible-playbook命令,可以使用Playbook来自动化执行任务和操作。例如,可以使用以下命令运行上述示例Playbook:

ansible-playbook -i inventory.ini  playbook.yml

常用参数

位置参数:
  playbook              Playbook文件名

选项:
  --flush-cache         清除清单中每个主机的事实缓存
  --force-handlers      即使任务失败也运行handlers
  --list-hosts          输出匹配的主机列表,不执行其他任何操作
  --list-tags           列出所有可用的标签
  --list-tasks          列出将要执行的所有任务
  --skip-tags SKIP_TAGS
                        只运行标签与这些值不匹配的plays和tasks
  --start-at-task START_AT_TASK
                        从与此名称匹配的任务开始运行playbook
  --step                逐步确认每个任务是否运行
  --syntax-check        对playbook执行语法检查,但不执行它
  -M MODULE_PATH, --module-path MODULE_PATH
                        在模块库之前添加以冒号分隔的路径(默认值为{{ ANSIBLE_HOME ~ "/plugins/modules:/usr/share/ansible/plugins/modules" })
  -e EXTRA_VARS, --extra-vars EXTRA_VARS
                        设置附加的变量,格式为key=value或YAML/JSON,如果是文件名则在前面加上@
  -f FORKS, --forks FORKS
                        指定要使用的并行进程数(默认值为5)
  -h, --help            显示此帮助消息并退出
  -i INVENTORY, --inventory INVENTORY, --inventory-file INVENTORY
                        指定清单主机路径或逗号分隔的主机列表。--inventory-file已弃用
  -k, --ask-pass        提示输入连接密码
  -l SUBSET, --limit SUBSET
                        进一步将选定的主机限制为附加的模式
  -t TAGS, --tags TAGS  只运行带有这些标签的plays和tasks
  -v, --verbose         导致Ansible打印更多的调试信息。多次添加-v会增加详细程度,内置插件目前支持最多-vvvvvv。开始时一个合理的级别是-vvv,连接调试可能需要-vvvv。

连接选项:
  --private-key PRIVATE_KEY_FILE, --key-file PRIVATE_KEY_FILE
                        使用此文件进行连接身份验证
  --scp-extra-args SCP_EXTRA_ARGS
                        指定传递给scp的额外参数(例如-l)
  --sftp-extra-args SFTP_EXTRA_ARGS
                        指定传递给sftp的额外参数(例如-f,-l)
  --ssh-common-args SSH_COMMON_ARGS
                        指定传递给sftp/scp/ssh的常见参数(例如ProxyCommand)
  --ssh-extra-args SSH_EXTRA_ARGS
                        指定传递给ssh的额外参数(例如-R)
  -T TIMEOUT, --timeout TIMEOUT
                        覆盖连接超时时间(以秒为单位)(默认值为10)
  -c CONNECTION, --connection CONNECTION
                        使用的连接类型(默认为smart)
  -u REMOTE_USER, --user REMOTE_USER
                        以此用户身份连接(默认为None)

特权升级选项:
  --become-method BECOME_METHOD
                        使用的特权升级方法(默认为sudo),使用`ansible-doc -t become -l`列出有效选择。
  --become-user BECOME_USER
                        以此用户身份运行操作(默认为root)
  -b, --become          以特权身份运行操作(不会提示输入密码)
阅读全文

Git使用GPG签名

生成秘钥

gpg --full-generate-key

列入所有keys

gpg --list-keys

输出如下,获取KeyID

pub   rsa4096 2023-06-12 [SC]
      <KeyID>
uid           [ultimate] name (comment) <email>
sub   rsa4096 2023-06-12 [E]

导出公钥,可以放到Github或者Gitlab。

gpg --armor --export <KeyID>

输入如下

-----BEGIN PGP PUBLIC KEY BLOCK-----
...
-----END PGP PUBLIC KEY BLOCK-----

设置Git

设置git commit时自动签名。

git config --global user.signingkey <KeyID>
git config --global gpg.program gpg
git config --global commit.gpgsign true

备份、恢复公私钥

以下是备份 GPG 公私钥的步骤:

打开命令行终端并输入以下命令来导出您的私钥:

gpg --export-secret-keys -a <KeyID> > private.key

输入以下命令以导出您的公钥:

gpg --export -a <KeyID> > public.key

将生成的两个文件(private.key 和 public.key)复制到您想要备份它们的位置,如 USB 磁盘或云存储。

您可以在需要恢复密钥时使用这些备份文件来重新导入密钥。使用以下命令将私钥导入到 GPG 中:

gpg --import private.key

最后,使用以下命令将公钥导入到 GPG 中:

gpg --import public.key

这样就完成了 GPG 公私钥的备份和恢复。

阅读全文

扩展Kubernetes(一)

Controller and CustomResourceDefinitions

Kubernetes的Custom Resource Definition(CRD)和Controller机制是Kubernetes提供的一种扩展Kubernetes API的方法,可以用于自定义Kubernetes资源类型以及对这些资源的控制。

CRD是一种自定义的Kubernetes资源类型,它允许用户在Kubernetes中创建新的资源类型。CRD通过Kubernetes API Server暴露出来,使得用户可以使用kubectl或其他工具对其进行管理。

Controller是CRD的另一个重要组成部分,它实现了对CRD资源的控制逻辑。Controller会根据CRD资源的状态变化来触发相应的操作,比如创建、更新、删除等。Controller通常会通过watch机制来监听CRD资源的变化,并根据实际情况对其进行处理。CRD和Controller机制为Kubernetes提供了更强大的扩展能力,使得用户可以更好地适应不同场景的需求。

Extend Kubernetes with CustomResourceDefinitions.

  • 定义、创建自定义资源
# crd.yaml
apiVersion: apiextensions.k8s.io/v1beta1
kind: CustomResourceDefinition
metadata:
  name: foos.samplecontroller.k8s.io
spec:
  group: samplecontroller.k8s.io
  version: v1alpha1
  names:
    kind: Foo
    plural: foos
  scope: Namespaced

kubectl apply -f crd.yaml
  • Kubectl操作自定义资源
kubectl get foo
kubectl watch foo
  • Client-go操作自定义资源
import (
    "context"

    samplecontrollerv1alpha1 "k8s.io/sample-controller/pkg/apis/samplecontroller/v1alpha1"
    clientset "k8s.io/sample-controller/pkg/generated/clientset/versioned"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)

func main() {
    // 获取 kubernetes 集群的 rest.Config 对象
    config, err := rest.InClusterConfig()
    if err != nil {
        panic(err)
    }
    
    // 创建 CRD 客户端
    cs, err := clientset.NewForConfig(config)
    if err != nil {
        panic(err)
    }

    // 创建自定义对象 foo
    foo := &samplecontrollerv1alpha1.Foo{
        ObjectMeta: metav1.ObjectMeta{
            Name:      "example-foo",
            Namespace: "default",
        },
        Spec: samplecontrollerv1alpha1.FooSpec{
            Size: 3,
        },
    }

    // 创建或更新自定义对象 foo
    result, err := cs.SamplecontrollerV1alpha1().Foos("default").Create(context.Background(), foo, metav1.CreateOptions{})
    if err != nil {
        result, err = cs.SamplecontrollerV1alpha1().Foos("default").Update(context.Background(), foo, metav1.UpdateOptions{})
        if err != nil {
            panic(err)
        }
    }

    // 删除自定义对象 foo
    err = cs.SamplecontrollerV1alpha1().Foos("default").Delete(context.Background(), foo.Name, metav1.DeleteOptions{})
    if err != nil {
        panic(err)
    }
}

Write Controller with client-go

Controller Architecture

CRD Client code generate

k8s.io/code-generator 是Kubernetes官方提供的一种代码生成工具。这个工具可以根据定义好的CRD来自动生成客户端代码,包括类型定义、列表和单个对象的获取、创建、更新和删除等操作的接口,以及相关实现代码。

通过使用k8s.io/code-generator, 我们可以避免手动编写大量重复的代码,从而提高开发效率和代码质量。

SCRIPT_ROOT=$(dirname "${BASH_SOURCE[0]}")/..
CODEGEN_PKG=${CODEGEN_PKG:-$(cd "${SCRIPT_ROOT}"; ls -d -1 ./vendor/k8s.io/code-generator 2>/dev/null || echo ../code-generator)}

bash "${CODEGEN_PKG}"/generate-groups.sh "deepcopy,client,informer,lister" \
k8s.io/sample-controller/pkg/generated k8s.io/sample-controller/pkg/apis \
samplecontroller:v1alpha1 \
--output-base "$(dirname "${BASH_SOURCE[0]}")/../../.." \
--go-header-file "${SCRIPT_ROOT}"/hack/boilerplate.go.txt
# pkg/generated
├── clientset
│   └── versioned
│       ├── clientset.go
│       ├── doc.go
│       ├── fake
│       │   ├── clientset_generated.go
│       │   ├── doc.go
│       │   └── register.go
│       ├── scheme
│       │   ├── doc.go
│       │   └── register.go
│       └── typed
│           └── samplecontroller
│               └── v1alpha1
│                   ├── doc.go
│                   ├── fake
│                   │   ├── doc.go
│                   │   ├── fake_foo.go
│                   │   └── fake_samplecontroller_client.go
│                   ├── foo.go
│                   ├── generated_expansion.go
│                   └── samplecontroller_client.go
├── informers
│   └── externalversions
│       ├── factory.go
│       ├── generic.go
│       ├── internalinterfaces
│       │   └── factory_interfaces.go
│       └── samplecontroller
│           ├── interface.go
│           └── v1alpha1
│               ├── foo.go
│               └── interface.go
└── listers
    └── samplecontroller
        └── v1alpha1
            ├── expansion_generated.go
            └── foo.go

Client-go components

  • Reflector: 定义在type Reflector inside package cache,
    通过Kubernetes API监控特定资源,资源可以为kubernetes内建资源也可以是自定义资源.
    当reflector收到来自API的资源更新或创建通知,它将通过相应的API创建一个对象,并将它推入Delta Fifo队列.

  • Informer: 定义在base controller inside package cache,
    它从Delta Fifo队列弹出一个对象保存下来,同时调用我们的controller并传递该对象.

  • Indexer: 定义在type Indexer inside package cache.
    为对象提供索引功能,一个典型的使用场景是基于对象的labels创建索引.
    Indexer维护索引通过一系列的索引函数,并使用一个线程安全的池子保存对象和它们的key.

Custom Controller components

  • Informer reference: Informer实例引用,
    需要我们在我们的controller代码中自己创建.

  • Indexer reference: Indexer实例引用,
    需要我们在我们的controller代码中自己创建,我们在执行处理过程中通过它来取回对象.

  • Resource Event Handlers: 一系列调函数,通过它们Informer传递对象给我们的controller.
    一个典型的模式为该回调函数获得对象的key后将对象的key推入workqueue为接下来的处理.

  • Work queue: 解耦对象传递和对象处理过程的队列.

  • Process Item: 对象具体的处理过程: 一般会使用Indexer reference或者Listing包装函数通过对象key取回对象.

import (
    "context"
    
    samplecontrollerv1alpha1 "k8s.io/sample-controller/pkg/apis/samplecontroller/v1alpha1"
    clientset "k8s.io/sample-controller/pkg/generated/clientset/versioned"
    informers "k8s.io/sample-controller/pkg/generated/informers/externalversions/samplecontroller/v1alpha1"
    listers "k8s.io/sample-controller/pkg/generated/listers/samplecontroller/v1alpha1"

    "k8s.io/apimachinery/pkg/runtime"
    "k8s.io/apimachinery/pkg/util/wait"
    "k8s.io/client-go/tools/cache"
    "k8s.io/client-go/util/workqueue"
)

type Controller struct {
    // 客户端,用于操作CRD
    clientset clientset.Interface
    
    // Informer 机制,确保 CRD 数据的及时获取和同步到内存中
    fooInformer informers.FooInformer
    fooLister   listers.FooLister
    
    // 工作队列,用于控制协程数目和并发度
    workqueue workqueue.RateLimitingInterface
    
    // 自定义对象转换接口
    scheme *runtime.Scheme
}

func NewController(clientset clientset.Interface, fooInformer informers.FooInformer) *Controller {
    c := &Controller{
        clientset: clientset,
        fooInformer: fooInformer,
        fooLister: fooInformer.Lister(),
        workqueue: workqueue.NewNamedRateLimitingQueue(workqueue.DefaultControllerRateLimiter(), "Foos"),
        scheme: runtime.NewScheme(),
    }
    // 注册自定义 API 对象到 scheme 中
    samplecontrollerv1alpha1.AddToScheme(c.scheme)
    // 绑定各种事件触发的回调函数
    fooInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: c.enqueueFoo,
        UpdateFunc: func(old, new interface{}) {
            c.enqueueFoo(new)
        },
        DeleteFunc: c.enqueueFooForDelete,
    })
    return c
}

func (c *Controller) Run(stopCh <-chan struct{}) error {
    defer c.workqueue.ShutDown()

    if !cache.WaitForCacheSync(stopCh, c.fooInformer.Informer().HasSynced) {
        return fmt.Errorf("failed to wait for caches to sync")
    }

    // 启动多个协程,处理工作队列中待处理的任务
    for i := 0; i < numWorkers; i++ {
        go wait.Until(c.runWorker, time.Second, stopCh)
    }

    <-stopCh
    return nil
}

func (c *Controller) runWorker() {
    for c.processNextWorkItem() {
    }
}

func (c *Controller) processNextWorkItem() bool {
    obj, shutdown := c.workqueue.Get()
    if shutdown {
        return false
    }
    defer c.workqueue.Done(obj)

    key := obj.(string)
    namespace, name, err := cache.SplitMetaNamespaceKey(key)
    if err != nil {
        c.workqueue.Forget(obj)
        return true
    }

    foo, err := c.fooLister.Foos(namespace).Get(name)
    if err != nil {
        if errors.IsNotFound(err) {
            c.workqueue.Forget(obj)
            return true
        }
        c.workqueue.AddRateLimited(obj)
        return true
    }

    // TODO: 根据业务逻辑处理自定义对象 foo
    
    // 处理完一个任务,将其从队列中删除
    c.workqueue.Forget(obj)
    return true
}

func (c *Controller) enqueueFoo(obj interface{}) {
    foo := obj.(*samplecontrollerv1alpha1.Foo)
    key, err := cache.MetaNamespaceKeyFunc(foo)
    if err != nil {
        return
    }
    c.workqueue.Add(key)
}

func (c *Controller) enqueueFooForDelete(obj interface{}) {
    foo, ok := obj.(*samplecontrollerv1alpha1.Foo)
    if !ok {
        tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
        if !ok {
            return
        }
        foo, ok = tombstone.Obj.(*samplecontrollerv1alpha1.Foo)
        if !ok {
            return
        }
    }
    key, err := cache.MetaNamespaceKeyFunc(foo)
    if err == nil {
        c.workqueue.Add(key)
    }
}

Controller Runtime

Controller-runtime是Kubernetes官方提供的一种Go语言工具包,用于编写自定义控制器。该工具包基于Kubernetes API机制,封装了常见的控制器开发模式,如Watch Reconcile、Leader Election、Event Emitting等。

Controller-runtime为Kubernetes控制器的开发提供了很多便利,使得开发者可以更加专注于核心业务逻辑的实现。

Controller-runtime Architecture

  1. Manager:管理器

Manager 是整个框架的核心,其负责以下几个任务:

  • 创建并启动各种控制器(Controller)
  • 启动 Web 服务器,提供健康检查和指标监控等服务
  • 提供客户端(Client)实例,用于操作 kubernetes API
  • 可以通过配置文件或环境变量来自定义控制器运行参数
  1. Client:客户端

Client 是操作 kubernetes API 的客户端,通常情况下我们使用 kubernetes/client-go 库即可,但是在 controller-runtime 框架中,它对原生的 client-go 库进行了封装,提供了更加易用的接口,同时还支持对多个版本的 API 对象进行操作。

  1. Cache:缓存

Cache 是用于缓存 kubernetes API 对象的组件,它可以高效地从 kubernetes API 中获取对象,并将其同步到内存中,以便控制器快速地获取和处理对象。

  1. Controller:控制器

Controller 是控制器的核心组件,用于监听和处理 kubernetes API 对象的变化,并在需要时调用 Reconciler 进行协调。

  1. Reconciler:协调器

Reconciler 是一个接口,用于协调处理 kubernetes API 对象的状态。在 controller-runtime 中,每个控制器都必须关联一个 Reconciler 实现,以便在对象出现变化时调用该 Reconciler 进行处理。Reconciler 接口定义如下:

type Reconciler interface {
    // 计算出与期望状态不同的部分
    Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error)
}

在 Reconcile 函数中,我们可以根据业务逻辑计算出当前对象的期望状态,然后与 kubernetes API 中的实际状态进行比较,如果两者不同,则更新实际状态,否则直接返回。

Cotroller-runtime 提供了一套完整的架构,封装了大量的功能代码,让开发 controller 变得更加简单方便。在使用 controller-runtime 构建 controller 时,我们继续只需要关注 Reconciler接口的实现即可。

Write Controller with controller-runtime

import (
    "context"

    samplecontrollerv1alpha1 "k8s.io/sample-controller/pkg/apis/samplecontroller/v1alpha1"
    "k8s.io/apimachinery/pkg/runtime"
    "k8s.io/apimachinery/pkg/util/wait"
    "k8s.io/apimachinery/pkg/api/errors"

    ctrl "sigs.k8s.io/controller-runtime"
    "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)

func main() {
    // 获取 kubernetes 集群的 rest.Config 对象
    config, err := rest.InClusterConfig()
    if err != nil {
        panic(err)
    }

    // 创建 kubernetes API 对象的编解码器
    scheme := runtime.NewScheme()
    samplecontrollerv1alpha1.AddToScheme(scheme)

    // 创建 manager 实例
    mgr, err := ctrl.NewManager(config, ctrl.Options{
        Scheme: scheme,
    })
    if err != nil {
        panic(err)
    }

    // 创建 Reconciler 实现
    reconciler := &FooReconciler{
        client: mgr.GetClient(),
        scheme: scheme,
    }

    // 创建控制器对象并注册到manager
    ctrl.NewControllerManagedBy(mgr).
        For(&samplecontrollerv1alpha1.Foo{}).
        WithEventFilter(predicate.GenerationChangedPredicate{}).
        Complete(reconciler)

    // 启动 manager,开始运行 controller
    if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
        panic(err)
    }
}

// FooReconciler 实现了 Reconciler 接口
type FooReconciler struct {
    client client.Client
    scheme *runtime.Scheme
}

func (r *FooReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
    // 获取 kubernetes API 对象
    foo := &samplecontrollerv1alpha1.Foo{}
    if err := r.client.Get(ctx, req.NamespacedName, foo); err != nil {
        if errors.IsNotFound(err) {
            return ctrl.Result{}, nil
        }
        return ctrl.Result{}, err
    }

    // 更新或创建 kubernetes API 对象
    if err := controllerutil.SetControllerReference(foo, foo, r.scheme); err != nil {
        return ctrl.Result{}, err
    }
    _, err := r.client.CreateOrUpdate(ctx, foo)
    if err != nil {
        return ctrl.Result{}, err
    }

    return ctrl.Result{}, nil
}
阅读全文

Gunicorn源码解析

深入浅出Python开发Web中已经实现了一个简单的WSGI服务器,但是在生产应用需要考虑健壮、性能、并发等等问题。

Gunicorn是一个Python的WSGI HTTP服务器实现,它支持多个进程模型,可以自动管理工作进程。在Python web开发中被广泛应用,是最受欢迎的服务器之一,许多著名的Python web框架都支持使用Gunicorn作为服务器。

Start

$ pip install gunicorn
$ cat myapp.py
  def app(environ, start_response):
      data = b"Hello, World!\n"
      start_response("200 OK", [
          ("Content-Type", "text/plain"),
          ("Content-Length", str(len(data)))
      ])
      return iter([data])
$ gunicorn -w 4 myapp:app
[2014-09-10 10:22:28 +0000] [30869] [INFO] Listening at: http://127.0.0.1:8000 (30869)
[2014-09-10 10:22:28 +0000] [30869] [INFO] Using worker: sync
[2014-09-10 10:22:28 +0000] [30874] [INFO] Booting worker with pid: 30874
[2014-09-10 10:22:28 +0000] [30875] [INFO] Booting worker with pid: 30875

gunicorn通过pip安装提供了一个cli入口,可以查看setup.py看到入口函数gunicorn.app.wsgiapp:run

# https://github.com/benoitc/gunicorn/blob/20.1.0/setup.py#L114
entry_points="""
[console_scripts]
gunicorn=gunicorn.app.wsgiapp:run


[paste.server_runner]
main=gunicorn.app.pasterapp:serve
"""

# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/app/wsgiapp.py#L61
def run():
    """\
    The ``gunicorn`` command line runner for launching Gunicorn with
    generic WSGI applications.
    """
    from gunicorn.app.wsgiapp import WSGIApplication
    WSGIApplication("%(prog)s [OPTIONS] [APP_MODULE]").run()

实例化了WSGIApplication,然后执行run方法启动。

WSGIApplication

# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/app/base.py#L70

class BaseApplication(object):
    ...
    def run(self):
        try:
            Arbiter(self).run()
        except RuntimeError as e:
            print("\nError: %s\n" % e, file=sys.stderr)
            sys.stderr.flush()
            sys.exit(1)

WSGIApplication继承BaseApplicationrun实现,最终实例化Arbiter执行其run

Arbiter

Arbiter是gunicorn的核心类,负责着整个服务器的运行管理,run方法是服务器启动的核心入口。

# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/arbiter.py#L196

def run(self):
    "Main master loop."
    # 1.初始化实例,并信号监听以及sock
    self.start()
    util._setproctitle("master [%s]" % self.proc_name)

    try:
        # 2.创建worker直到worker数量满足用户的指定条件,如果当前worker数量超过用户指定的条件,则会杀掉创建时间比较长的worker
        self.manage_workers()

        while True:
            # 3.负责判断该进程是否是正真的master,如果是则提升为正真的master(这一块放在最后一部分进行分析)
            self.maybe_promote_master()

            # 4.获取信号
            sig = self.SIG_QUEUE.pop(0) if self.SIG_QUEUE else None
            if sig is None:
                # 5.利用select休眠1小时
                self.sleep()
                # 6.判断worker是否超时,如果是则杀掉worker(将在Worker章节进行分析) 
                self.murder_workers()
                self.manage_workers()
                continue

            if sig not in self.SIG_NAMES:
                self.log.info("Ignoring unknown signal: %s", sig)
                continue

            signame = self.SIG_NAMES.get(sig)
            # 调用对应的信号处理
            handler = getattr(self, "handle_%s" % signame, None)
            if not handler:
                self.log.error("Unhandled signal: %s", signame)
                continue
            self.log.info("Handling signal: %s", signame)
            handler()
            # 7.这样下次循环就不会等待一秒了
            self.wakeup()
    # 8.服务异常,进行退出处理
    except (StopIteration, KeyboardInterrupt):
        # 收到用户的退出信号(按下CTRL+C) 
        self.halt()
    except HaltServer as inst:
        # Worker运行异常的时候
        self.halt(reason=inst.reason, exit_status=inst.exit_status)
    except SystemExit:
        raise
    except Exception:
        # 其它的运行异常
        self.log.info("Unhandled exception in main loop",
                      exc_info=True)
        self.stop(False)
        if self.pidfile is not None:
            self.pidfile.unlink()
        sys.exit(-1)

run方法的loop中会执行核心的manage_workers管理子进程逻辑。

ManageWorkers

# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/arbiter.py#L545

class Arbiter(object):

    def manage_workers(self):
        if len(self.WORKERS) < self.num_workers:
            self.spawn_workers()
        ...

    def spawn_worker(self):
        self.worker_age += 1
        worker = self.worker_class(self.worker_age, self.pid, self.LISTENERS,
                                   self.app, self.timeout / 2.0,
                                   self.cfg, self.log) # 通过配置获取子进程worker类

        self.cfg.pre_fork(self, worker)
        pid = os.fork()      # 执行fork
        if pid != 0:
            worker.pid = pid
            self.WORKERS[pid] = worker
            return pid

        ...

        # Process Child
        worker.pid = os.getpid()
        try:
            util._setproctitle("worker [%s]" % self.proc_name)
            self.log.info("Booting worker with pid: %s", worker.pid)
            self.cfg.post_fork(self, worker)
            worker.init_process()        # 执行worker的初始化
            sys.exit(0)
        except SystemExit:
            raise


    def spawn_workers(self):

        for _ in range(self.num_workers - len(self.WORKERS)):
            self.spawn_worker()
            time.sleep(0.1 * random.random())

manage_workers中会实例化配置的worker类,然后fork子进程执行workerinit_process

Worker

Worker在gunicorn的pre-fork子进程中负责运行应用处理请求, 大部分worker都是对gunicorn的http、wsgi接口封装,但可以通过自定义协议封装来支持TCP应用。

  • 如果worker使用gunicorn实现的http、wsgi接口,则只需要考虑怎么处理来自客户端的请求即可。Gunicron内置实现了基于thread、gevent、tornado和sync的不同worker,默认gunicorn使用sync的worker,处理请求是单线程、阻塞同步模式。
  • 如果worker在init_process之后,worker不再使用gunicorn实现的http、wsgi接口,仅仅使用共享master进程获得的socket,则可以自定义协议实现TCP应用,这种模式下gunicorn充当一个TCP应用进程管理器。
# src/gunicorn/gunicorn/workers git:(master) $ ls -lh

-rw-rw-r-- 1 arvin arvin 5.6K Jul 11  2022 base_async.py
-rw-rw-r-- 1 arvin arvin 9.0K Mar 24 10:50 base.py
-rw-rw-r-- 1 arvin arvin 6.0K Mar 24 10:50 geventlet.py
-rw-rw-r-- 1 arvin arvin 5.7K Mar 24 10:50 ggevent.py
-rw-rw-r-- 1 arvin arvin  12K Jul 11  2022 gthread.py
-rw-rw-r-- 1 arvin arvin 5.9K Mar 24 10:50 gtornado.py
-rw-rw-r-- 1 arvin arvin  594 Jul 11  2022 __init__.py
-rw-rw-r-- 1 arvin arvin 7.2K Jul 11  2022 sync.py
-rw-rw-r-- 1 arvin arvin 1.7K Jul 11  2022 workertmp.py

SyncWorker

# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/workers/base.py#L29
class Worker(object):

    def init_process(self):
        ....
        self.load_wsgi() # 加载wsgi应用
        ...
        self.run()


# https://github.com/benoitc/gunicorn/blob/20.1.0/gunicorn/workers/sync.py#L25
class SyncWorker(base.Worker):

    def accept(self, listener):
        client, addr = listener.accept()
        client.setblocking(1)
        util.close_on_exec(client)
        self.handle(listener, client, addr) # 处理连接请求

    def run_for_one(self, timeout):
        while self.alive:
            ...
           self.accept(listener)
            ...


    def run_for_multiple(self, timeout):
        while self.alive:
            ...
           self.accept(listener)
            ...

    def run(self): # main loop
        ...
        if len(self.sockets) > 1:
            self.run_for_multiple(timeout) # 进入请求处理loop
        else:
            self.run_for_one(timeout)  # 进入请求处理loop

    def handle(self, listener, client, addr):
         ...
         parser = http.RequestParser(self.cfg, client, addr) # 解析http请求
         req = next(parser)
         self.handle_request(listener, req, client, addr) # 处理wsgi请求
        ...


    def handle_request(self, listener, req, client, addr):
          ...
          resp, environ = wsgi.create(req, client, addr,
                                        listener.getsockname(), self.cfg) # 构建wsgi协议对象
          ...
          respiter = self.wsgi(environ, resp.start_response) # 调用wsgi应用
          ...
  • SyncWorker启动流程为:init_process(初始化)->self.load_wsgi(加载wsgi应用)->run(进入处理请求loop)。
  • SyncWorker处理一个http请求的流程为:accept(接受请求)->handle(解析http)->handle_request(封装wsgi)->self.wsgi(调用wsgi应用)。
阅读全文

深入浅出Python开发Web

深入浅出Web开发中聊了业务系统需要通过应用网关协议与Web服务器交互。

Browser <—HTTP协议—> Web Server <—应用网关协议—> 业务系统(PHP,Python,Java...)

WSGI

WSGI(Web Server Gateway Interface)是Python Web应用程序与Web服务器之间的通信协议,它定义了Web服务器如何与Web框架进行交互,从而实现了Web服务器和Web框架之间的解耦。

WSGI的作用就像是一个桥梁,它把Web服务器和Web框架连接起来,让它们能够进行有效的通信。Web服务器只需要遵循WSGI协议,将请求和响应传递给WSGI应用程序,而不需要了解具体的应用程序实现细节。同样地,Web框架也只需编写符合WSGI规范的应用程序,而不需要考虑与特定的Web服务器进行交互的问题。这种解耦的方式使得Web开发更加灵活和可扩展。

WSGI协议的实现非常简单。它只要求应用程序提供一个callable对象,接受两个参数:一个环境变量字典和一个可调用的start_response函数。Web服务器调用该callable,并传递环境变量字典和start_response函数作为参数。应用程序可以使用环境变量字典获取请求信息,然后通过调用start_response函数返回响应头信息。应用程序还需要返回一个可迭代的响应body对象。

# Web Server 进程
...
# 加载wsgi应用
...
# 处理Request调用
def wsgi_callable(environ, start_response):
    ...
    return body
...
# 返回Response给客户端
...

因为WSGI协议的简单实现,任何Python Web框架都可以轻松地与任何符合WSGI协议的Web服务器进行交互,从而实现灵活和可扩展的Web应用程序开发。

目前社区的Python Web框架,比如Flask、Django、Madara最终都是暴露一个WSGI的Callable对象作为应用的入口。

SimpleWSGIServer

详解HTTP协议中实现的SimpleHTTPServer改造成一个支持WSGI协议的Web服务器。

import socket  # 导入 socket 模块
from multiprocessing.dummy import Pool as ThreadPool  # 导入线程池模块ThreadPool,用于多线程处理请求
import io  # 导入io模块,用于发送和接收HTTP报文
import traceback  # 导入traceback模块,用于打印错误信息
import logging  # 导入logging模块,用于记录日志



class Server(object):
    # 定义一个类变量 SERVER_STRING,存储服务器的名称和版本信息
    SERVER_STRING = b"Server: SimpleHttpd/1.0.0\r\n"


    def __init__(self, host, port, worker_count=4):
        self._host = host  # 存储主机名
        self._port = port  # 存储端口号
        self._listen_fd = None  # 存储监听套接字
        self._worker_count = worker_count  # 存储工作线程数
        self._worker_pool = ThreadPool(worker_count)  # 创建指定数量的工作线程池
        self._logger = logging.getLogger("simple.httpd")  # 创建日志记录器
        self._logger.setLevel(logging.DEBUG)  # 设置日志级别为DEBUG
        self._logger.addHandler(logging.StreamHandler())  # 添加输出日志到控制台的处理器


    def run(self):
        self._listen_fd = socket.socket(socket.AF_INET, socket.SOCK_STREAM)  # 创建TCP/IP套接字
        self._listen_fd.bind((self._host, self._port))  # 绑定IP地址和端口号
        self._listen_fd.listen(self._worker_count)  # 监听客户端连接请求
        try:
            while True:
                conn, addr = self._listen_fd.accept()  # 接受客户端连接并返回新的套接字和地址
                self._worker_pool.apply_async(self.accept_request, (conn, addr,))  # 将连接套接字和地址交给工作线程异步执行处理
        except Exception as e:
            traceback.print_exc()  # 打印异常堆栈信息
        finally:
            self._listen_fd.close()  # 关闭监听套接字


    def accept_request(self, conn: socket.socket, addr):
        try:
            method, path, http_version, req_headers, req_body = self.recv_request(conn)
            status = ""
            # 构建environ字典
            environ = {
                'REQUEST_METHOD': method,  # 请求方法,例如 'POST'
                'SCRIPT_NAME': '',  # 脚本名称
                'PATH_INFO': path,  # 路径信息,例如'/hello/world'
                'QUERY_STRING': '',  # 查询字符串,例如'a=1&b=2'
                'CONTENT_TYPE': '',  # 请求体的类型,例如'application/json'
                'CONTENT_LENGTH': len(req_body),  # 请求体的长度(单位为字节)
                'SERVER_NAME': self._host,  # 服务器主机名
                'SERVER_PORT': str(self._port),  # 服务器端口号
                'HTTP_HOST': req_headers.get('Host', ''),
                'HTTP_USER_AGENT': req_headers.get('User-Agent', ''),  # HTTP请求头中的'User-Agent'字段
                'HTTP_ACCEPT': req_headers.get('Accept', ''),  # HTTP请求头中的'Accept'字段
                'HTTP_ACCEPT_LANGUAGE': req_headers.get('Accept-Language', ''),  # HTTP请求头中的'Accept-Language'字段
                'HTTP_ACCEPT_ENCODING': req_headers.get('Accept-Encoding', ''),  # HTTP请求头中的'Accept-Encoding'字段
                'HTTP_CONNECTION': req_headers.get('Connection', '')  # HTTP请求头中的'Connection'字段
            }
            environ['wsgi.input'] = io.BytesIO(req_body)
            # 执行wsgi应用
            body = self.wsgi_application(environ, self.build_start_response(conn, status))
            # 迭代响应body对象,返回给客户端
            for bt in body:
                self.send_response(conn, None, bt, None)
            self._logger.info("{}:{} {} {} {} {}".format(addr[0], addr[1], http_version, method, path, status))  # 记录日志信息
        except Exception as e:
            traceback.print_exc()  # 打印异常堆栈信息
        finally:
            conn.close()  # 关闭连接套接字


    def wsgi_application(self, environ, start_response):
        """
        wsgi应用
        """
        data = b"Hello, World!\n"
        start_response("200 OK", [
            ("Content-Type", "text/plain"),
            ("Content-Length", str(len(data)))
        ])
        return [data]


    def build_start_response(self, conn, the_status):
        """
        构建start_response函数
        """
        def start_response(status, headers):
            the_status = status
            self.send_response(conn, status=status, headers=headers)
        return start_response


    def send_response(self, conn: socket.socket, status=None, body=None, headers=None):
        if not status is None:
            conn.sendall("HTTP/1.0 {}\r\n".format(status).encode())
            conn.sendall(self.SERVER_STRING)
        if not headers is None:
            for header in headers:
                conn.sendall("{}: {}\r\n".format(*header).encode())
        if not body is None:
            if not isinstance(body, bytes):
                body = str(body).encode()
            conn.sendall(b"\r\n")
            conn.sendall(body)


    def recv_request(self, conn: socket.socket):
        # 读取请求行
        line = b''
        while not line.endswith(b'\r\n'):
            data = conn.recv(1)
            if not data:
                raise ConnectionError('Connection closed unexpectedly')
            line += data
        method, path, version = line.strip().decode().split(' ', 2)


        # 读取请求头
        headers = {}
        while True:
            line = b''
            while not line.endswith(b'\r\n'):
                data = conn.recv(1)
                if not data:
                    raise ConnectionError('Connection closed unexpectedly')
                line += data
            if line == b'\r\n':
                break
            key, value = line.strip().decode().split(': ', 1)
            headers[key] = value


        # 读取请求体
        content_length = int(headers.get('Content-Length', '0'))
        if content_length > 0:
            body = conn.recv(content_length)
        else:
            body = b""


        # 返回请求行、请求头和请求体
        return method, path, version, headers, body



if __name__ == "__main__":
    server = Server("0.0.0.0", 3000)  # 创建服务器实例
    server.run()  # 启动服务器
阅读全文