欢迎光临
我们一直在努力

Flink SQL多租户实践:Yarn/K8s环境下的资源隔离

1. 多租户架构设计

1.1 多租户核心概念

资源隔离与共享的平衡策略。

多租户架构层次
├── 物理隔离层
│ ├── 独立集群(完全隔离)
│ ├── 专用节点(节点级隔离)
│ └── 网络分区(网络级隔离)
├── 逻辑隔离层
│ ├── 命名空间(Namespace隔离)
│ ├── 资源配额(Quota限制)
│ └── 权限控制(RBAC授权)
└── 运行时隔离层
├── 容器隔离(cgroups/namespace)
├── 资源限制(CPU/Memory限制)
└── 网络策略(NetworkPolicy)

1.2 租户模型设计

基于业务场景的租户分类。

# 租户定义模板
apiVersion: flink.apache.org/v1beta1
kind: Tenant
metadata:
name: datascienceteam
labels:
department: datascience
environment: production
priority: high
spec:
# 资源配额
resourceQuota:
cpu: "100" # 100核
memory: "200Gi" # 200GB内存
gpu: "4" # 4张GPU卡
storage: "1Ti" # 1TB存储

# 网络隔离
networkPolicy:
ingress:
from:
namespaceSelector:
matchLabels:
tenant: datascienceteam
ports:
protocol: TCP
port: 8081 # Flink Web UI
egress:
to:
ipBlock:
cidr: "10.0.0.0/8"
ports:
protocol: TCP
port: 9092 # Kafka

# 存储隔离
storageClasses:
name: fastssd
accessModes: [ "ReadWriteOnce" ]
storage: "100Gi"
name: archivehdd
accessModes: [ "ReadWriteMany" ]
storage: "500Gi"

# 调度策略
scheduling:
nodeSelector:
tenant: datascience
tolerations:
key: "gpu"
operator: "Equal"
value: "true"
effect: "NoSchedule"
affinity:
podAntiAffinity:
preferredDuringSchedulingIgnoredDuringExecution:
weight: 100
podAffinityTerm:
labelSelector:
matchLabels:
app: flinktaskmanager
topologyKey: "kubernetes.io/hostname"

2. YARN环境多租户实践

2.1 YARN队列资源配置

基于Capacity Scheduler的多租户管理。

<!– capacity-scheduler.xml – YARN队列配置 –>
<configuration>
<!– 根队列配置 –>
<property>
<name>yarn.scheduler.capacity.root.queues</name>
<value>prod,dev,research,batch</value>
</property>

<!– 生产队列 – 高优先级 –>
<property>
<name>yarn.scheduler.capacity.root.prod.capacity</name>
<value>50</value> <!– 50%资源 –>
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.maximum-capacity</name>
<value>80</value> <!– 最大80% –>
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.minimum-user-limit-percent</name>
<value>25</value> <!– 单用户最少25% –>
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.state</name>
<value>RUNNING</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.acl_submit_applications</name>
<value>prod_team</value> <!– 提交权限 –>
</property>

<!– 开发队列 – 中等优先级 –>
<property>
<name>yarn.scheduler.capacity.root.dev.capacity</name>
<value>30</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.dev.user-limit-factor</name>
<value>2</value> <!– 可超卖2倍 –>
</property>

<!– 研究队列 – 低优先级,可抢占 –>
<property>
<name>yarn.scheduler.capacity.root.research.capacity</name>
<value>15</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.research.preemptable</name>
<value>true</value> <!– 允许被抢占 –>
</property>

<!– 批处理队列 – 最低优先级 –>
<property>
<name>yarn.scheduler.capacity.root.batch.capacity</name>
<value>5</value>
</property>

<!– 全局配置 –>
<property>
<name>yarn.scheduler.capacity.maximum-applications</name>
<value>10000</value> <!– 最大应用数 –>
</property>
<property>
<name>yarn.scheduler.capacity.maximum-am-resource-percent</name>
<value>0.2</value> <!– AM最多20%资源 –>
</property>
</configuration>

2.2 Flink on YARN多租户配置

针对不同队列的Flink作业配置。

#!/bin/bash
# multi-tenant-flink-job.sh

# 公共配置
FLINK_HOME="/opt/flink"
YARN_QUEUE="$1" # 队列参数
JOB_NAME="$2"
JAR_PATH="$3"

# 根据队列设置不同资源配置
case $YARN_QUEUE in
"prod")
# 生产队列 – 高资源保障
CONTAINER_MEMORY="8192" # 8GB
TASKMANAGER_MEMORY="6144" # 6GB
JOBMANAGER_MEMORY="2048" # 2GB
SLOTS_PER_TM="4"
PARALLELISM="16"
;;
"dev")
# 开发队列 – 中等资源
CONTAINER_MEMORY="4096" # 4GB
TASKMANAGER_MEMORY="3072" # 3GB
JOBMANAGER_MEMORY="1024" # 1GB
SLOTS_PER_TM="2"
PARALLELISM="8"
;;
"research")
# 研究队列 – 低资源,可抢占
CONTAINER_MEMORY="2048" # 2GB
TASKMANAGER_MEMORY="1536" # 1.5GB
JOBMANAGER_MEMORY="512" # 0.5GB
SLOTS_PER_TM="1"
PARALLELISM="4"
;;
*)
echo "Unknown queue: $YARN_QUEUE"
exit 1
;;
esac

# 提交Flink作业到指定队列
$FLINK_HOME/bin/flink run \\
-m yarn-cluster \\
-yqu $YARN_QUEUE \\
-ynm "$JOB_NAME" \\
-yD taskmanager.memory.process.size=${TASKMANAGER_MEMORY}m \\
-yD jobmanager.memory.process.size=${JOBMANAGER_MEMORY}m \\
-yD taskmanager.numberOfTaskSlots=$SLOTS_PER_TM \\
-yD parallelism.default=$PARALLELISM \\
-yD high-availability=zookeeper \\
-yD high-availability.storageDir=hdfs:///flink/ha/$YARN_QUEUE \\
-yD state.checkpoints.dir=hdfs:///flink/checkpoints/$YARN_QUEUE \\
-yD yarn.application.name="flink-$YARN_QUEUE$JOB_NAME" \\
-c com.example.StreamingJob \\
$JAR_PATH

# 使用示例
# ./multi-tenant-flink-job.sh prod order-processing-job /jobs/order-processing.jar
# ./multi-tenant-flink-job.sh dev data-exploration-job /jobs/exploration.jar

2.3 动态资源分配策略

基于负载的动态资源调整。

— 资源使用监控表
CREATE TABLE resource_usage_metrics (
queue_name STRING,
job_id STRING,
user_name STRING,
allocated_memory_gb DOUBLE,
used_memory_gb DOUBLE,
allocated_vcores INT,
used_vcores INT,
queue_capacity_percent DOUBLE,
measurement_time TIMESTAMP(3),
WATERMARK FOR measurement_time AS measurement_time INTERVAL '30' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'yarn-metrics',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);

— 自动伸缩策略
CREATE TABLE auto_scaling_decisions (
decision_time TIMESTAMP(3),
queue_name STRING,
action STRING, — SCALE_UP, SCALE_DOWN, REBALANCE
current_capacity DOUBLE,
new_capacity DOUBLE,
reason STRING,
priority INT
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:postgresql://scheduler-db:5432/decisions',
'table-name' = 'scaling_decisions'
);

— 基于负载的动态队列容量调整
INSERT INTO auto_scaling_decisions
SELECT
CURRENT_TIMESTAMP AS decision_time,
queue_name,
CASE
WHEN avg_usage > 0.8 AND queue_capacity < maximum_capacity THEN 'SCALE_UP'
WHEN avg_usage < 0.3 AND queue_capacity > minimum_capacity THEN 'SCALE_DOWN'
ELSE 'MAINTAIN'
END AS action,
queue_capacity AS current_capacity,
CASE
WHEN avg_usage > 0.8 THEN LEAST(queue_capacity * 1.2, maximum_capacity)
WHEN avg_usage < 0.3 THEN GREATEST(queue_capacity * 0.8, minimum_capacity)
ELSE queue_capacity
END AS new_capacity,
'Automatic scaling based on usage: ' || CAST(avg_usage AS STRING) AS reason,
CASE
WHEN queue_name = 'prod' THEN 1 — 生产队列高优先级
WHEN queue_name = 'dev' THEN 2
ELSE 3
END AS priority
FROM (
SELECT
queue_name,
AVG(used_memory_gb / allocated_memory_gb) AS avg_usage,
AVG(queue_capacity_percent) AS queue_capacity,
80.0 AS maximum_capacity, — 队列最大容量
20.0 AS minimum_capacity — 队列最小容量
FROM resource_usage_metrics
WHERE measurement_time > CURRENT_TIMESTAMP INTERVAL '10' MINUTE
GROUP BY queue_name, TUMBLE(measurement_time, INTERVAL '5' MINUTE)
);

3. Kubernetes多租户实践

3.1 命名空间隔离配置

基于Namespace的租户隔离。

# namespace-quotas.yaml – 命名空间资源配额
apiVersion: v1
kind: Namespace
metadata:
name: flinkprod
labels:
tenant: production
environment: prod
flink-version: "1.16"

apiVersion: v1
kind: ResourceQuota
metadata:
name: flinkprodquota
namespace: flinkprod
spec:
hard:
# 计算资源限制
requests.cpu: "200" # 200核
requests.memory: 400Gi # 400GB内存
limits.cpu: "400" # 400核上限
limits.memory: 800Gi # 800GB内存上限

# 对象数量限制
pods: "100" # 最大Pod数
services: "50" # 最大Service数
configmaps: "100" # 最大ConfigMap数
persistentvolumeclaims: "50" # 最大PVC数
secrets: "100" # 最大Secret数

# 存储限制
requests.storage: "5Ti" # 5TB存储
persistentvolumeclaims: "50"

apiVersion: v1
kind: LimitRange
metadata:
name: flinkprodlimits
namespace: flinkprod
spec:
limits:
type: Container
default: # 默认资源限制
cpu: "2" # 默认2核
memory: "4Gi" # 默认4GB内存
defaultRequest: # 默认资源请求
cpu: "1" # 最少1核
memory: "2Gi" # 最少2GB内存
max: # 最大资源限制
cpu: "16" # 单容器最大16核
memory: "32Gi" # 单容器最大32GB内存
min: # 最小资源限制
cpu: "100m" # 单容器最少0.1核
memory: "512Mi" # 单容器最少512MB内存

3.2 Flink Kubernetes Operator多租户配置

使用Operator管理多租户Flink集群。

# flink-cluster-multi-tenant.yaml
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: orderprocessingprod
namespace: flinkprod
labels:
tenant: production
application: orderprocessing
priority: high
spec:
image: flink:1.16.0
flinkVersion: v1_16
imagePullPolicy: IfNotPresent

# 服务账户和权限
serviceAccount: flinkprodsa

# 资源配额
resourceQuota:
requests:
cpu: "100"
memory: "200Gi"
limits:
cpu: "200"
memory: "400Gi"

# JobManager配置
jobManager:
resource:
memory: "2048m"
cpu: 2
replicas: 1
podTemplate:
spec:
nodeSelector:
node-type: highmemory
tolerations:
key: "dedicated"
operator: "Equal"
value: "flink-prod"
effect: "NoSchedule"

# TaskManager配置
taskManager:
resource:
memory: "4096m"
cpu: 4
replicas: 10
podTemplate:
spec:
nodeSelector:
node-type: highmemory
tolerations:
key: "dedicated"
operator: "Equal"
value: "flink-prod"
effect: "NoSchedule"
affinity:
podAntiAffinity:
preferredDuringSchedulingIgnoredDuringExecution:
weight: 100
podAffinityTerm:
labelSelector:
matchLabels:
app: flinktaskmanager
topologyKey: "kubernetes.io/hostname"

# Pod安全上下文
podSecurityContext:
runAsUser: 1000
runAsGroup: 1000
fsGroup: 1000

# Flink配置
flinkConfiguration:
taskmanager.numberOfTaskSlots: "2"
parallelism.default: "20"
high-availability: org.apache.flink.kubernetes.highavailability.KubernetesHaServicesFactory
high-availability.storageDir: file:///flink/ha/
state.backend: filesystem
state.checkpoints.dir: file:///flink/checkpoints/
state.savepoints.dir: file:///flink/savepoints/
web.cancel.enable: "false" # 生产环境禁用Web取消

# 日志配置
logging:
log4j-console.properties: |
rootLogger.level = INFO
rootLogger.appenderRef.console.ref = ConsoleAppender
logger.jobmanager.name = org.apache.flink.runtime.jobmaster.JobMaster
logger.jobmanager.level = INFO

# 网络策略
networkPolicy:
ingress:
from:
namespaceSelector:
matchLabels:
tenant: monitoring
ports:
port: 9249 # Prometheus metrics
from:
namespaceSelector:
matchLabels:
tenant: dataplatform
ports:
port: 8081 # Flink Web UI

3.3 网络策略隔离

基于NetworkPolicy的网络隔离。

# network-policies.yaml – 网络隔离策略
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: flinkprodisolation
namespace: flinkprod
spec:
podSelector:
matchLabels:
app: flink
policyTypes:
Ingress
Egress

# 入站规则
ingress:
# 允许监控命名空间访问指标端口
from:
namespaceSelector:
matchLabels:
tenant: monitoring
ports:
protocol: TCP
port: 9249 # Prometheus
protocol: TCP
port: 8081 # Flink Web UI

# 允许数据平台命名空间访问
from:
namespaceSelector:
matchLabels:
tenant: dataplatform
ports:
protocol: TCP
port: 8081
protocol: TCP
port: 6123 # TaskManager数据端口

# 出站规则
egress:
# 允许访问Kafka
to:
namespaceSelector:
matchLabels:
tenant: kafka
ports:
protocol: TCP
port: 9092

# 允许访问HDFS
to:
ipBlock:
cidr: 10.0.100.0/24 # HDFS集群IP段
ports:
protocol: TCP
port: 8020 # HDFS NameNode
protocol: TCP
port: 50010 # HDFS DataNode

# 允许访问数据库
to:
namespaceSelector:
matchLabels:
tenant: database
ports:
protocol: TCP
port: 5432 # PostgreSQL
protocol: TCP
port: 3306 # MySQL

# 允许DNS查询
to:
namespaceSelector: {}
ports:
protocol: UDP
port: 53

4. 存储多租户隔离

4.1 持久化存储隔离

租户级存储配置。

# storage-classes.yaml – 存储类配置
apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
name: fastssdprod
labels:
tenant: production
performance-tier: high
provisioner: kubernetes.io/gcepd
parameters:
type: pdssd
replication-type: regionalpd
reclaimPolicy: Retain # 保留数据
allowVolumeExpansion: true
volumeBindingMode: WaitForFirstConsumer

apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
name: standardhdddev
labels:
tenant: development
performance-tier: standard
provisioner: kubernetes.io/gcepd
parameters:
type: pdstandard
reclaimPolicy: Delete # 开发环境可删除
allowVolumeExpansion: true
volumeBindingMode: Immediate

4.2 租户存储配额

存储资源限制与隔离。

# pvc-templates.yaml – 租户PVC模板
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: flinkcheckpointsprod
namespace: flinkprod
annotations:
tenant: production
usage: checkpoints
spec:
accessModes:
ReadWriteMany
resources:
requests:
storage: 1Ti # 1TB存储
storageClassName: fastssdprod
selector:
matchLabels:
tenant: production

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: flinkhaprod
namespace: flinkprod
annotations:
tenant: production
usage: highavailability
spec:
accessModes:
ReadWriteOnce
resources:
requests:
storage: 100Gi # 100GB HA存储
storageClassName: fastssdprod

5. 安全与权限控制

5.1 RBAC权限隔离

基于角色的访问控制。

# rbac-multi-tenant.yaml
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: flinktenantadmin
rules:
apiGroups: ["flink.apache.org"]
resources: ["flinkdeployments", "flinksessions"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
apiGroups: [""]
resources: ["pods", "services", "configmaps", "persistentvolumeclaims"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]

apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
namespace: flinkprod
name: flinkprodoperator
rules:
apiGroups: [""]
resources: ["pods", "services", "configmaps", "secrets"]
verbs: ["get", "list", "watch", "create", "update", "patch", "delete"]
apiGroups: [""]
resources: ["pods/exec"]
verbs: ["create"]
apiGroups: ["networking.k8s.io"]
resources: ["networkpolicies"]
verbs: ["get", "list", "watch"]

apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: flinkprodbinding
namespace: flinkprod
subjects:
kind: ServiceAccount
name: flinkprodsa
namespace: flinkprod
roleRef:
kind: Role
name: flinkprodoperator
apiGroup: rbac.authorization.k8s.io

5.2 服务账户隔离

租户级服务账户配置。

# service-accounts.yaml
apiVersion: v1
kind: ServiceAccount
metadata:
name: flinkprodsa
namespace: flinkprod
labels:
tenant: production
app: flink
secrets:
name: flinkprodtoken

apiVersion: v1
kind: Secret
metadata:
name: flinkprodtoken
namespace: flinkprod
annotations:
kubernetes.io/service-account.name: flinkprodsa
type: kubernetes.io/serviceaccounttoken

6. 监控与成本核算

6.1 多租户监控

租户级资源使用监控。

— 租户资源使用统计
CREATE TABLE tenant_resource_usage (
tenant_id STRING,
namespace STRING,
resource_type STRING, — cpu, memory, storage, gpu
allocated_amount DOUBLE,
used_amount DOUBLE,
utilization_rate DOUBLE,
cost_usd DOUBLE,
measurement_time TIMESTAMP(3),
PRIMARY KEY (tenant_id, resource_type, measurement_time) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:postgresql://cost-db:5432/tenant-billing',
'table-name' = 'resource_usage'
);

— 资源使用率监控
CREATE VIEW tenant_utilization_alerts AS
SELECT
tenant_id,
namespace,
resource_type,
utilization_rate,
measurement_time,
CASE
WHEN utilization_rate < 0.3 THEN 'UNDER_UTILIZED'
WHEN utilization_rate > 0.8 THEN 'OVER_UTILIZED'
ELSE 'HEALTHY'
END AS utilization_status
FROM tenant_resource_usage
WHERE utilization_rate < 0.3 OR utilization_rate > 0.8;

— 成本分摊计算
CREATE TABLE tenant_cost_allocation (
tenant_id STRING,
period_start DATE,
period_end DATE,
total_cost DOUBLE,
cpu_cost DOUBLE,
memory_cost DOUBLE,
storage_cost DOUBLE,
network_cost DOUBLE,
allocated_by STRING — fair-share, actual-usage, reserved
) WITH (
'connector' = 'jdbc',
'table-name' = 'cost_allocation'
);

INSERT INTO tenant_cost_allocation
SELECT
tenant_id,
DATE_TRUNC('MONTH', measurement_time) AS period_start,
DATE_TRUNC('MONTH', measurement_time) + INTERVAL '1' MONTH INTERVAL '1' DAY AS period_end,
SUM(cost_usd) AS total_cost,
SUM(CASE WHEN resource_type = 'cpu' THEN cost_usd ELSE 0 END) AS cpu_cost,
SUM(CASE WHEN resource_type = 'memory' THEN cost_usd ELSE 0 END) AS memory_cost,
SUM(CASE WHEN resource_type = 'storage' THEN cost_usd ELSE 0 END) AS storage_cost,
SUM(CASE WHEN resource_type = 'network' THEN cost_usd ELSE 0 END) AS network_cost,
'actual-usage' AS allocated_by
FROM tenant_resource_usage
WHERE measurement_time >= CURRENT_TIMESTAMP INTERVAL '1' MONTH
GROUP BY tenant_id, DATE_TRUNC('MONTH', measurement_time);

6.2 自动伸缩与优化

基于使用率的自动资源调整。

— 自动伸缩推荐
CREATE TABLE scaling_recommendations (
recommendation_id STRING,
tenant_id STRING,
resource_type STRING,
current_allocation DOUBLE,
recommended_allocation DOUBLE,
estimated_savings DOUBLE,
confidence_score DOUBLE,
recommendation_type STRING, — scale-up, scale-down, optimize
reason STRING,
generated_at TIMESTAMP(3)
) WITH ('connector' = 'jdbc');

INSERT INTO scaling_recommendations
SELECT
MD5(CONCAT(tenant_id, resource_type, CAST(measurement_time AS STRING))) AS recommendation_id,
tenant_id,
resource_type,
allocated_amount AS current_allocation,
CASE
WHEN utilization_rate < 0.3 THEN allocated_amount * 0.8 — 缩容20%
WHEN utilization_rate > 0.8 THEN allocated_amount * 1.2 — 扩容20%
ELSE allocated_amount
END AS recommended_allocation,
CASE
WHEN utilization_rate < 0.3 THEN (allocated_amount * 0.2) * cost_per_unit
ELSE 0
END AS estimated_savings,
CASE
WHEN ABS(utilization_rate 0.6) < 0.1 THEN 0.9 — 接近最优,高置信度
ELSE 0.7
END AS confidence_score,
CASE
WHEN utilization_rate < 0.3 THEN 'scale-down'
WHEN utilization_rate > 0.8 THEN 'scale-up'
ELSE 'optimize'
END AS recommendation_type,
'Utilization rate: ' || CAST(utilization_rate AS STRING) AS reason,
CURRENT_TIMESTAMP AS generated_at
FROM tenant_resource_usage
WHERE utilization_rate < 0.3 OR utilization_rate > 0.8;

7. 总结

多租户环境下的资源隔离需要多层次、细粒度的控制策略:

核心隔离维度
资源隔离:CPU、内存、存储的硬性限制
网络隔离:NetworkPolicy实现网络层隔离
存储隔离:StorageClass和PVC的租户绑定
权限隔离:RBAC和ServiceAccount的权限控制
最佳实践
分层配额:根据租户优先级设置不同资源配额
弹性伸缩:基于使用率的自动资源调整
成本透明:详细的资源使用计量和成本分摊
安全加固:最小权限原则和网络策略
运维建议
监控告警:实时监控租户资源使用情况
容量规划:基于历史数据的容量预测
自动化运维:自动化的资源调整和优化
文档化:清晰的租户SLA和资源使用规范
通过完善的多租户隔离方案,可以在共享集群中实现资源公平性、安全隔离性和运维效率的平衡。

赞(0)
未经允许不得转载:171主机测评 » Flink SQL多租户实践:Yarn/K8s环境下的资源隔离
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址