|
| 1 | +apiVersion: config.karmada.io/v1alpha1 |
| 2 | +kind: ResourceInterpreterCustomization |
| 3 | +metadata: |
| 4 | + name: declarative-configuration-flinkdeployment |
| 5 | +spec: |
| 6 | + target: |
| 7 | + apiVersion: flink.apache.org/v1beta1 |
| 8 | + kind: FlinkDeployment |
| 9 | + customizations: |
| 10 | + healthInterpretation: |
| 11 | + luaScript: > |
| 12 | + function InterpretHealth(observedObj) |
| 13 | + if observedObj.status ~= nil and observedObj.status.jobStatus ~= nil then |
| 14 | + return observedObj.status.jobStatus.state ~= 'CREATED' and observedObj.status.jobStatus.state ~= 'RECONCILING' |
| 15 | + end |
| 16 | + return false |
| 17 | + end |
| 18 | + replicaResource: |
| 19 | + luaScript: > |
| 20 | + local kube = require("kube") |
| 21 | +
|
| 22 | + local function isempty(s) |
| 23 | + return s == nil or s == '' |
| 24 | + end |
| 25 | +
|
| 26 | + function GetReplicas(observedObj) |
| 27 | + -- FlinkDeployments presently will not be subdivided among clusters, replica should be 1 |
| 28 | + replica = 1 |
| 29 | + requires = { |
| 30 | + resourceRequest = {}, |
| 31 | + } |
| 32 | + -- Add jobmanager resources into replica requirement |
| 33 | +
|
| 34 | + jm_replicas = observedObj.spec.jobManager.replicas |
| 35 | + if isempty(jm_replicas) then |
| 36 | + jm_replicas = 1 |
| 37 | + end |
| 38 | +
|
| 39 | + for i = 1, jm_replicas do |
| 40 | + requires.resourceRequest.cpu = kube.resourceAdd(requires.resourceRequest.cpu, tostring(observedObj.spec.jobManager.resource.cpu)) |
| 41 | + requires.resourceRequest.memory = kube.resourceAdd(requires.resourceRequest.memory, observedObj.spec.jobManager.resource.memory) |
| 42 | + end |
| 43 | +
|
| 44 | + -- Add task manager resources into replica requirement |
| 45 | +
|
| 46 | + parallelism = observedObj.spec.job.parallelism |
| 47 | + tms = math.ceil(parallelism / observedObj.spec.flinkConfiguration['taskmanager.numberOfTaskSlots']) |
| 48 | +
|
| 49 | + for i = 1, tms do |
| 50 | + requires.resourceRequest.cpu = kube.resourceAdd(requires.resourceRequest.cpu, tostring(observedObj.spec.taskManager.resource.cpu)) |
| 51 | + requires.resourceRequest.memory = kube.resourceAdd(requires.resourceRequest.memory, observedObj.spec.taskManager.resource.memory) |
| 52 | + end |
| 53 | +
|
| 54 | + return replica, requires |
| 55 | + end |
| 56 | + statusAggregation: |
| 57 | + luaScript: > |
| 58 | + function AggregateStatus(desiredObj, statusItems) |
| 59 | + if statusItems == nil then |
| 60 | + return desiredObj |
| 61 | + end |
| 62 | + if desiredObj.status == nil then |
| 63 | + desiredObj.status = {} |
| 64 | + end |
| 65 | + clusterInfo = {} |
| 66 | + jobManagerDeploymentStatus = '' |
| 67 | + jobStatus = {} |
| 68 | + lifecycleState = '' |
| 69 | + observedGeneration = 0 |
| 70 | + reconciliationStatus = {} |
| 71 | + taskManager = {} |
| 72 | +
|
| 73 | + for i = 1, #statusItems do |
| 74 | + currentStatus = statusItems[i].status |
| 75 | + if currentStatus ~= nil then |
| 76 | + clusterInfo = currentStatus.clusterInfo |
| 77 | + jobManagerDeploymentStatus = currentStatus.jobManagerDeploymentStatus |
| 78 | + jobStatus = currentStatus.jobStatus |
| 79 | + observedGeneration = currentStatus.observedGeneration |
| 80 | + lifecycleState = currentStatus.lifecycleState |
| 81 | + reconciliationStatus = currentStatus.reconciliationStatus |
| 82 | + taskManager = currentStatus.taskManager |
| 83 | + end |
| 84 | + end |
| 85 | +
|
| 86 | + desiredObj.status.clusterInfo = clusterInfo |
| 87 | + desiredObj.status.jobManagerDeploymentStatus = jobManagerDeploymentStatus |
| 88 | + desiredObj.status.jobStatus = jobStatus |
| 89 | + desiredObj.status.lifecycleState = lifecycleState |
| 90 | + desiredObj.status.observedGeneration = observedGeneration |
| 91 | + desiredObj.status.reconciliationStatus = reconciliationStatus |
| 92 | + desiredObj.status.taskManager = taskManager |
| 93 | + return desiredObj |
| 94 | + end |
| 95 | + statusReflection: |
| 96 | + luaScript: > |
| 97 | + function ReflectStatus(observedObj) |
| 98 | + status = {} |
| 99 | + if observedObj == nil or observedObj.status == nil then |
| 100 | + return status |
| 101 | + end |
| 102 | + status.clusterInfo = observedObj.status.clusterInfo |
| 103 | + status.jobManagerDeploymentStatus = observedObj.status.jobManagerDeploymentStatus |
| 104 | + status.jobStatus = observedObj.status.jobStatus |
| 105 | + status.observedGeneration = observedObj.status.observedGeneration |
| 106 | + status.lifecycleState = observedObj.status.lifecycleState |
| 107 | + status.reconciliationStatus = observedObj.status.reconciliationStatus |
| 108 | + status.taskManager = observedObj.status.taskManager |
| 109 | + return status |
| 110 | + end |
0 commit comments