ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink Native Kubernetes 部署实战:会话模式、应用模式与 Pod 模板配置全指南

Flink Native Kubernetes 部署实战:会话模式、应用模式与 Pod 模板配置全指南 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本指南基于 Apache Flink 的 Kubernetes 原生集成Native Kubernetes资源提供方完整讲解如何将 Flink 直接部署到运行中的 Kubernetes 集群并借助 Flink 与 Kubernetes APIServer 的直接通信实现 TaskManager 的动态伸缩。读完本文你将掌握 Session/Application 两种部署模式的启动命令与-D参数覆盖技巧、Web UI 的三种暴露方式、日志与插件管理、Secrets 挂载、HA 配置以及 Pod 模板的字段覆盖规则并能结合源码理解其底层实现。快速上手在 Kubernetes 上启动第一个 Flink 集群集成原理简介Kubernetes 是目前最流行的容器编排系统用于自动化应用的部署、扩缩容与管理。Flink 的 Native Kubernetes 集成允许你直接在一个运行中的 Kubernetes 集群上部署 FlinkFlink 客户端通过 Fabric8 Kubernetes client 与 Kubernetes APIServer 通信创建/删除 Deployment、Pod、ConfigMap、Service 等资源并 watch Pod 与 ConfigMapResourceManager 还能根据作业所需资源动态分配与释放 TaskManager。Apache Flink 同时提供 Kubernetes Operator 用于管理 Kubernetes 上的 Flink 集群它同时支持 standalone 与 native 两种部署模式可大幅简化 Flink 资源的部署、配置与生命周期管理详见 Flink Kubernetes Operator 文档。从源码结构看这套集成的核心实现位于 flink-kubernetes 模块其中 KubernetesClusterClientFactory.java 负责识别--target kubernetes-*目标并创建集群描述符KubernetesClusterDescriptor.java 则实现会话/应用/单作业三种集群的部署、回收与终止。前置准备开始前请确保你拥有一个满足以下条件的 Kubernetes 集群Kubernetes 1.9KubeConfig可通过~/.kube/config访问集群并具备对 pods、services 的 list、create、delete 权限可用kubectl auth can-i list|create|edit|delete pods校验启用 Kubernetes DNSdefault服务账号具备创建、删除 Pod 所需的 RBAC 权限。若搭建集群遇到困难可参考 Kubernetes 官方安装文档。启动一个 Flink Session 集群集群就绪且kubectl指向它之后即可按 Session 模式 启动 Flink 集群# (1) 启动 Kubernetes Session 集群 $ ./bin/kubernetes-session.sh -Dkubernetes.cluster-idmy-first-flink-cluster # (2) 提交示例作业 $ ./bin/flink run \ --target kubernetes-session \ -Dkubernetes.cluster-idmy-first-flink-cluster \ ./examples/streaming/TopSpeedWindowing.jar # (3) 通过删除集群 Deployment 停止 Session $ kubectl delete deployment/my-first-flink-cluster提示默认情况下 Flink 的 Web UI 与 REST 端点以ClusterIP类型 Service 暴露访问方式请参考下文 访问 Flink Web UI。Session 的两种运行模式Session 模式支持两种运行方式detached 模式默认kubernetes-session.sh将 Flink 集群部署到 Kubernetes 后即退出attached 模式-Dexecution.attachedtruekubernetes-session.sh保持运行允许输入命令控制正在运行的集群例如输入stop停止 Session 集群输入help查看所有支持的命令。重新挂接re-attach到 cluster-id 为my-first-flink-cluster的集群$ ./bin/kubernetes-session.sh \ -Dkubernetes.cluster-idmy-first-flink-cluster \ -Dexecution.attachedtrue停止运行中的 Session 集群可直接删除 Flink Deployment或通过管道输入stop命令$ echo stop | ./bin/kubernetes-session.sh \ -Dkubernetes.cluster-idmy-first-flink-cluster \ -Dexecution.attachedtruebin/flink与bin/kubernetes-session.sh均支持通过-Dkeyvalue覆盖 Flink 配置文件 中的任意配置项。部署模式Application 与 Session生产环境建议采用 Application 模式该模式能为应用提供更好的隔离性。两种模式的核心区别、类加载与用户 JAR 的识别规则如下Session 模式启动命令中指定的 JAR 被识别为用户 JARApplication 模式启动命令指定的 JAR 以及 Flinkusrlib目录下的所有 JAR 都会被加入用户类路径。Application 模式打包用户代码进镜像Application 模式要求用户代码与 Flink 镜像捆绑在一起因为该模式下用户代码的main()方法在集群上执行应用终止后所有 Flink 组件会被正确清理。捆绑可通过修改基础 Flink Docker 镜像完成也可借助 User Artifact Management 上传/下载本地不可用的工件。Flink 社区提供了基础 Docker 镜像用于打包用户代码FROM flink RUN mkdir -p $FLINK_HOME/usrlib COPY /path/of/my-flink-job.jar $FLINK_HOME/usrlib/my-flink-job.jar镜像构建并发布为custom-image-name后即可启动 Application 集群$ ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-idmy-first-application-cluster \ -Dkubernetes.container.image.refcustom-image-name \ local:///opt/flink/usrlib/my-flink-job.jarApplication 模式使用 User Artifact Management如果本机已有作业 JAR可启用工件上传Flink 会在部署期间把本地工件上传到 DFS再由 JobManager Pod 拉取$ ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-idmy-first-application-cluster \ -Dkubernetes.container.imagecustom-image-name \ -Dkubernetes.artifacts.local-upload-enabledtrue \ -Dkubernetes.artifacts.local-upload-targets3://my-bucket/ \ local:///tmp/my-flink-job.jarkubernetes.artifacts.local-upload-enabled对应源码 KubernetesConfigOptions.java 中的LOCAL_UPLOAD_ENABLED默认false启用该特性kubernetes.artifacts.local-upload-targetLOCAL_UPLOAD_TARGET必须指向一个存在且权限配置正确的远程目标另外可通过kubernetes.artifacts.local-upload-overwriteLOCAL_UPLOAD_OVERWRITE默认false控制是否覆盖远程已存在的工件。还可以通过user.artifacts.artifact-list追加混合了本地与远程的额外工件注意;需转义为\;$ ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-idmy-first-application-cluster \ -Dkubernetes.container.imagecustom-image-name \ -Dkubernetes.artifacts.local-upload-enabledtrue \ -Dkubernetes.artifacts.local-upload-targets3://my-bucket/ \ -Duser.artifacts.artifact-listlocal:///tmp/my-flink-udf1.jar\;s3://my-bucket/my-flink-udf2.jar \ local:///tmp/my-flink-job.jar警告本地上传不会覆盖已存在的远程工件若作业 JAR 或附加工件已通过 DFS 或 HTTP(S) 远程可达Flink 会在 JobManager Pod 上直接拉取# FileSystem $ ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-idmy-first-application-cluster \ -Dkubernetes.container.imagecustom-image-name \ s3://my-bucket/my-flink-job.jar # HTTP(S) $ ./bin/flink run-application \ --target kubernetes-application \ -Dkubernetes.cluster-idmy-first-application-cluster \ -Dkubernetes.container.imagecustom-image-name \ https://ip:port/my-flink-job.jar提示Application 模式的 JAR 拉取支持文件系统或 HTTP(S)。JAR 会被下载到镜像内的user.artifacts.base-dir/kubernetes.namespace/kubernetes.cluster-id路径下参见 user.artifacts.base-dir 与 kubernetes.namespace。Application 集群部署完成后可与之交互# 列出集群上运行的作业 $ ./bin/flink list --target kubernetes-application -Dkubernetes.cluster-idmy-first-application-cluster # 取消运行中的作业 $ ./bin/flink cancel --target kubernetes-application -Dkubernetes.cluster-idmy-first-application-cluster jobIdkubernetes.cluster-id用于指定集群名称且必须唯一未指定时 Flink 会自动生成随机名称。从源码看KubernetesClusterClientFactory.createClusterDescriptor 在配置缺失时会以flink-cluster-前缀加随机 ID 生成且截断到Constants.MAXIMUM_CHARACTERS_OF_CLUSTER_ID45 字符以内。kubernetes.container.image.ref用于指定启动 Pod 所用的镜像未配置时默认使用与当前 Flink/Scala 版本匹配的官方镜像见 CONTAINER_IMAGE。Kubernetes 上的 Flink 配置核心配置项Kubernetes 专属配置项的完整列表见配置页面。Flink 使用 Fabric8 Kubernetes client 与 Kubernetes APIServer 通信负责创建/删除 Deployment、Pod、ConfigMap、Service 等资源并 watch Pod 与 ConfigMap。除上述 Flink 配置外Fabric8 客户端的一些专家选项可通过系统属性或环境变量配置。例如通过 Flink 配置项设置并发最大请求数containerized.master.env.KUBERNETES_MAX_CONCURRENT_REQUESTS: 200 env.java.opts.jobmanager: -Dkubernetes.max.concurrent.requests200常用 Kubernetes 配置项速查默认值均取自 KubernetesConfigOptions.java配置项默认值说明kubernetes.cluster-id无自动生成集群唯一标识≤45 字符仅允许小写字母数字与-格式[a-z](https://link.gitcode.com/i/c01e75978ba131a102edfd8110a4514e)kubernetes.container.image.ref随版本变化的官方镜像启动 Pod 所用镜像须与应用使用的 Flink/Scala 版本一致kubernetes.namespacedefaultFlink 集群所在的命名空间kubernetes.service-accountdefaultJM/TM 使用的服务账号可被下面两个覆盖kubernetes.jobmanager.service-accountdefaultJobManager 创建 TaskManager Pod 时使用的服务账号kubernetes.taskmanager.service-accountdefaultTaskManager watch leader ConfigMap 时使用的服务账号kubernetes.rest-service.exposed.typeClusterIPREST 服务暴露类型kubernetes.jobmanager.replicas1同时启动的 JobManager Pod 数1 可启动备选 JobManagerkubernetes.jobmanager.cpu.amount1.0JobManager 使用的 CPU 数kubernetes.taskmanager.cpu.amount-1.0TaskManager CPU 数默认按每槽位 1 CPU 计算kubernetes.artifacts.local-upload-enabledfalse是否在部署前上传local://工件到 DFSkubernetes.artifacts.local-upload-overwritefalse是否覆盖远程已存在工件kubernetes.artifacts.local-upload-target无本地工件上传的远程 DFS 目录kubernetes.transactional-operation.max-retries15Kubernetes 事务性操作如checkAndUpdateConfigMap最大重试次数kubernetes.transactional-operation.initial-retry-delay50ms事务性操作重试的初始延迟kubernetes.transactional-operation.max-retry-delay1min事务性操作重试的最大延迟kubernetes.client.io-pool.size4Kubernetes 客户端执行阻塞 IO如启停 TM Pod、更新 leader ConfigMap的线程池大小kubernetes.hostnetwork.enabledfalse是否启用 HostNetwork 模式kubernetes.decorator.hadoop-conf-mount.enabledtrue是否挂载 Hadoop 配置使用 Flink Kubernetes Operator 且外部挂载时需设为falsekubernetes.decorator.kerberos-mount.enabledtrue是否挂载 Kerberos 配置外部挂载时需设为false访问 Flink 的 Web UI通过kubernetes.rest-service.exposed.type配置项Flink 的 Web UI 与 REST 端点可以多种方式暴露取值定义见 REST_SERVICE_EXPOSED_TYPEClusterIP默认在集群内部 IP 上暴露服务仅集群内可达。如需访问 JobManager UI 或向现有 Session 提交作业需要启动本地代理之后即可通过localhost:8081查看仪表盘或提交作业$ kubectl port-forward service/ServiceName 8081NodePort在每个节点的 IP 上以静态端口NodePort暴露服务可通过NodeIP:NodePort访问 JobManager 服务LoadBalancer借助云厂商负载均衡器对外暴露服务。由于云厂商与 Kubernetes 需要时间准备负载均衡器客户端日志中可能先出现NodePort形式的 JobManager Web 界面可用kubectl get services/cluster-id-rest获取 EXTERNAL-IP并手工构造http://EXTERNAL-IP:8081访问。更多信息请参考 Kubernetes 官方文档 publishing services in Kubernetes。警告视环境而定以LoadBalancer暴露 REST 服务可能使集群公开可访问通常意味着可执行任意代码请谨慎使用。日志、插件与自定义镜像日志管理Kubernetes 集成会把conf/log4j-console.properties与conf/logback-console.xml以 ConfigMap 形式暴露给 Pod修改这些文件后对新启动的集群生效。访问日志默认情况下 JobManager 与 TaskManager 同时把日志输出到控制台和每个 Pod 内的/opt/flink/logSTDOUT/STDERR只重定向到控制台。可通过以下命令查看$ kubectl logs pod-namePod 运行中时也可用kubectl exec -it pod-name bash进入容器查看日志或调试进程。查看 TaskManager 日志Flink 会自动释放空闲的 TaskManager 以避免资源浪费这会让对应 Pod 的日志难以获取。可通过配置 resourcemanager.taskmanager-timeout 延长空闲 TaskManager 被释放前的时间以便有更多时间检查日志文件。动态修改日志级别若日志器已配置为自动检测配置变更可通过修改 ConfigMap 动态调整日志级别假设 cluster-id 为my-first-flink-cluster$ kubectl edit cm flink-config-my-first-flink-cluster使用插件要使用插件必须把插件复制到 Flink JobManager/TaskManager Pod 中的正确位置。内置插件无需挂载卷或构建自定义镜像即可使用参见内置插件。例如为 Flink Session 集群启用 S3 插件$ ./bin/kubernetes-session.sh \ -Dcontainerized.master.env.ENABLE_BUILT_IN_PLUGINSflink-s3-fs-hadoop-version.jar \ -Dcontainerized.taskmanager.env.ENABLE_BUILT_IN_PLUGINSflink-s3-fs-hadoop-version.jarversion请替换为当前 Flink 版本号。自定义 Docker 镜像如需使用自定义镜像可通过kubernetes.container.image.ref指定。Flink 社区提供了功能丰富的 Flink Docker 镜像作为起点关于如何启用插件、添加依赖等请参考如何定制 Flink 的 Docker 镜像。使用 SecretsKubernetes Secrets 用于存放密码、令牌、密钥等少量敏感数据避免直接写入 Pod 规格或镜像。Flink on Kubernetes 支持两种使用方式将 Secrets 作为文件挂载下面的命令把 secretmysecret挂载到启动 Pod 的/path/to/secret路径下$ ./bin/kubernetes-session.sh -Dkubernetes.secretsmysecret:/path/to/secret此后mysecret的用户名与密码分别存放在/path/to/secret/username与/path/to/secret/password文件中。配置值格式为foo:/opt/secrets-foo,bar:/opt/secrets-bar见 KUBERNETES_SECRETS。详见 Kubernetes 官方文档 Using Secrets as files from a pod。将 Secrets 作为环境变量下面的命令把 secretmysecret暴露为启动 Pod 中的环境变量$ ./bin/kubernetes-session.sh -Dkubernetes.env.secretKeyRef\ env:SECRET_USERNAME,secret:mysecret,key:username;\ env:SECRET_PASSWORD,secret:mysecret,key:password环境变量SECRET_USERNAME与SECRET_PASSWORD分别包含mysecret的用户名与密码。配置值格式为env:FOO_ENV,secret:foo_secret,key:foo_key;env:BAR_ENV,secret:bar_secret,key:bar_key见 KUBERNETES_ENV_SECRET_KEY_REF。详见 Kubernetes 官方文档 Using Secrets as environment variables。高可用与资源管理Kubernetes 上的高可用Kubernetes 上的高可用可直接使用现有高可用服务。将 kubernetes.jobmanager.replicas 配置为大于 1可启动备选 JobManager 实现更快的恢复。注意启动备选 JobManager 时必须同时启用高可用参见 kubernetes.jobmanager.replicas 的源码说明。手动资源清理Flink 利用 Kubernetes 的 OwnerReference 垃圾回收机制 清理所有集群组件。Flink 创建的所有资源包括ConfigMap、Service、Pod的OwnerReference都被设置为deployment/cluster-id因此删除 Deployment 时所有相关资源会被自动删除$ kubectl delete deployment/cluster-id支持的 Kubernetes 版本与命名空间目前支持所有 1.9的 Kubernetes 版本。Kubernetes 的命名空间通过资源配额在多个用户间划分集群资源Flink on Kubernetes 可通过 kubernetes.namespace 指定启动集群的命名空间。RBAC 配置基于角色的访问控制RBAC按企业内用户角色调节对计算/网络资源的访问。用户可为 JobManager 配置访问 Kubernetes APIServer 所需的 RBAC 角色与服务账号。每个命名空间都有默认服务账号但default服务账号未必具备创建/删除 Pod 的权限需要更新其权限或指定绑定了合适角色的其他服务账号$ kubectl create clusterrolebinding flink-role-binding-default --clusterroleedit --serviceaccountdefault:default若不想使用default服务账号可创建新的flink-service-account并绑定角色再通过-Dkubernetes.service-accountflink-service-account配置 JobManager Pod 使用的服务账号用于创建/删除 TaskManager Pod 与 leader ConfigMap同时允许 TaskManager watch leader ConfigMap 以获取 JobManager 与 ResourceManager 的地址$ kubectl create serviceaccount flink-service-account $ kubectl create clusterrolebinding flink-role-binding-flink --clusterroleedit --serviceaccountdefault:flink-service-account从源码看JOB_MANAGER_SERVICE_ACCOUNT 与 TASK_MANAGER_SERVICE_ACCOUNT 默认均为default且可回退到kubernetes.service-account。更多信息参考 Kubernetes 官方 RBAC Authorization 文档。Pod 模板深度定制 JobManager 与 TaskManager PodFlink 允许通过模板文件定义 JobManager 与 TaskManager Pod以支持 Flink Kubernetes 配置项无法直接覆盖的高级特性。使用kubernetes.pod-template-file.default对应源码 KUBERNETES_POD_TEMPLATE是kubernetes.pod-template-file.jobmanager与kubernetes.pod-template-file.taskmanager的回退键指定包含 Pod 定义的本地文件用于初始化 JobManager 与 TaskManager。主容器必须命名为flink-main-container详见下文示例。被 Flink 覆盖的字段模板中部分字段会被 Flink 覆盖其有效值解析机制分三类由 Flink 定义Defined by Flink用户不可配置由用户定义Defined by the user用户可自由指定Flink 不会额外设置最终值遵循优先级显式配置项 Pod 模板值 配置项默认值与 Flink 合并Merged with FlinkFlink 会把用户定义值与自身值合并同名字段以 Flink 值为准。下表完整列出会被覆盖的 Pod 字段未列出的模板字段不受影响。Pod Metadata字段类别相关配置项说明name由 Flink 定义—JobManager Pod 名被覆盖为 kubernetes.cluster-id 对应的 Deployment 名TaskManager Pod 名被覆盖为 Flink ResourceManager 生成的clusterID-attempt-index模式namespace由用户定义kubernetes.namespaceJobManager Deployment 与 TaskManager Pod 均创建在用户指定的命名空间ownerReferences由 Flink 定义—JM/TM Pod 的 owner reference 始终指向 JobManager Deployment可通过 kubernetes.jobmanager.owner.reference 控制 Deployment 的删除时机annotations由用户定义kubernetes.jobmanager.annotations、kubernetes.taskmanager.annotationsFlink 追加配置项指定的额外注解labels与 Flink 合并kubernetes.jobmanager.labels、kubernetes.taskmanager.labelsFlink 在用户定义值之上追加内部标签Pod Spec字段类别相关配置项说明imagePullSecrets由用户定义kubernetes.container.image.pull-secretsFlink 追加配置项指定的额外拉取密钥nodeSelector由用户定义kubernetes.jobmanager.node-selector、kubernetes.taskmanager.node-selectorFlink 追加配置项指定的额外节点选择器tolerations由用户定义kubernetes.jobmanager.tolerations、kubernetes.taskmanager.tolerationsFlink 追加配置项指定的额外容忍度restartPolicy由 Flink 定义—JobManager Pod 为always由 Deployment 保证重启TaskManager Pod 为never不应重启serviceAccount由用户定义kubernetes.service-accountJM/TM Pod 使用用户指定的服务账号创建volumes与 Flink 合并—Flink 追加内部 ConfigMap 卷如 flink-config-volume、hadoop-config-volume用于分发 Flink 与 Hadoop 配置Main Container Spec字段类别相关配置项说明env与 Flink 合并containerized.master.env.{ENV_NAME}、containerized.taskmanager.env.{ENV_NAME}Flink 在用户定义值之上追加内部环境变量image由用户定义kubernetes.container.image.ref镜像按用户定义值的优先级顺序解析imagePullPolicy由用户定义kubernetes.container.image.pull-policy镜像拉取策略按用户定义值的优先级顺序解析name由 Flink 定义—容器名被覆盖为flink-main-containerresources由用户定义内存jobmanager.memory.process.size、taskmanager.memory.process.sizeCPUkubernetes.jobmanager.cpu、kubernetes.taskmanager.cpu内存与 CPU含 requests/limits被 Flink 配置项覆盖其他资源如 ephemeral-storage保留containerPorts与 Flink 合并—Flink 追加内部容器端口rest、jobmanager-rpc、blob、taskmanager-rpcvolumeMounts与 Flink 合并—Flink 追加内部卷挂载如 flink-config-volume、hadoop-config-volume用于分发 Flink 与 Hadoop 配置Pod 模板示例pod-template.yaml展示了 initContainer 拉取用户 JAR、sidecar 容器收集日志、挂载 hostPath/emptyDir 卷等高级用法apiVersion: v1 kind: Pod metadata: name: jobmanager-pod-template spec: initContainers: - name: artifacts-fetcher image: busybox:latest # 用 wget 或其他工具从远端存储获取用户 JAR command: [ wget, https://path/of/StateMachineExample.jar, -O, /flink-artifact/myjob.jar ] volumeMounts: - mountPath: /flink-artifact name: flink-artifact containers: # 不要修改主容器名称 - name: flink-main-container resources: requests: ephemeral-storage: 2048Mi limits: ephemeral-storage: 2048Mi volumeMounts: - mountPath: /opt/flink/volumes/hostpath name: flink-volume-hostpath - mountPath: /opt/flink/artifacts name: flink-artifact - mountPath: /opt/flink/log name: flink-logs # 使用 sidecar 容器把日志推送到远端存储或做其他调试 - name: sidecar-log-collector image: sidecar-log-collector:latest command: [ command-to-upload, /remote/path/of/flink-logs/ ] volumeMounts: - mountPath: /flink-logs name: flink-logs volumes: - name: flink-volume-hostpath hostPath: path: /tmp type: Directory - name: flink-artifact emptyDir: { } - name: flink-logs emptyDir: { }用户 JAR 与类路径在 Kubernetes 上原生部署 Flink 时以下 JAR 会被识别为用户 JAR 并加入用户类路径Session 模式启动命令中指定的 JARApplication 模式启动命令中指定的 JAR 以及 Flinkusrlib目录下的所有 JAR。关于类加载的详细机制请参考调试类加载文档。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink CDC 在 Kubernetes 上的部署实战Session 模式与 Kubernetes Operator 模式完全指南Flink CDC 在 Kubernetes 上的部署实战Session 模式与 Kubernetes Operator 模式完全指南 Flink CDC 作后端数据集成大数据流处理变更数据捕获数据同步Apache Flink Native Kubernetes 部署实战指南从 Session 集群到 Application 模式Apache Flink Native Kubernetes 部署实战指南从 Session 集群到 Application 模式 本文基于当前仓库中 Fli大数据流处理批处理数据工程Flink CDC 部署模式全指南Standalone、YARN 与 Kubernetes 的安装、配置与作业提交Flink CDC 部署模式全指南Standalone、YARN 与 Kubernetes 的安装、配置与作业提交 本文是 Flink CDC 部署模块 d后端数据集成大数据流处理变更数据捕获数据同步创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表