|
@@ -43,6 +43,7 @@ class MyFlink:
|
|
|
job_ids = self.list_job_id()
|
|
|
for job_id in job_ids:
|
|
|
self.stop_by_id(job_id)
|
|
|
+ print(f"停止job:{job_id}")
|
|
|
|
|
|
# 根据job id 停止任务
|
|
|
def stop_by_id(self, job_id):
|
|
@@ -97,7 +98,7 @@ class MyTest(unittest.TestCase):
|
|
|
flink = MyFlink("https://flink.rxdptest.k5.bigtree.tech")
|
|
|
jar_id = flink.upload_jar("/Users/alvin/bigtree/rxdp-dm-jobs/rxdp-dm-jobs-flink/target/rxdp-dm-jobs-flink-1.0-SNAPSHOT.jar")
|
|
|
print(f"jar_id: {jar_id}")
|
|
|
- jar_id = "b33b6c52-0e68-4f11-b6af-aeb7ffb116bf_rxdp-dm-jobs-flink-1.0-SNAPSHOT.jar"
|
|
|
+ flink.stop_all()
|
|
|
print(flink.start_job(jar_id, """
|
|
|
--kafkaServer kafka-0.kafka-headless.rxdptest.svc.k5.bigtree.zone:9092,kafka-1.kafka-headless.rxdptest.svc.k5.bigtree.zone:9092,kafka-2.kafka-headless.rxdptest.svc.k5.bigtree.zone:9092
|
|
|
--esServer elasticsearch-master.rxdptest.svc.k5.bigtree.zone
|