feat(managed-job): add args and parallelism param to run flink job method
This commit is contained in:
parent
07b8a36e63
commit
5bc047dbd1
@ -36,6 +36,10 @@ spec:
|
||||
type: integer
|
||||
jarUri:
|
||||
type: string
|
||||
args:
|
||||
type: array
|
||||
items:
|
||||
type: string
|
||||
savepointInterval:
|
||||
type: string
|
||||
format: duration
|
||||
|
||||
@ -16,6 +16,7 @@ type FlinkJobSpec struct {
|
||||
JarURI string `json:"jarUri"`
|
||||
SavepointInterval metaV1.Duration `json:"savepointInterval"`
|
||||
EntryClass string `json:"entryClass"`
|
||||
Args []string `json:"args"`
|
||||
}
|
||||
|
||||
type FlinkJobStatus struct {
|
||||
|
||||
@ -42,6 +42,8 @@ func (job *ManagedJob) run(restoreMode bool) error {
|
||||
AllowNonRestoredState: true,
|
||||
EntryClass: job.def.Spec.EntryClass,
|
||||
SavepointPath: savepointPath,
|
||||
Parallelism: job.def.Spec.Parallelism,
|
||||
ProgramArg: job.def.Spec.Args,
|
||||
})
|
||||
if err == nil {
|
||||
pkg.Logger.Info("[managed-job] [run] jar successfully ran", zap.Any("run-jar-resp", runJarResp))
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user