|
22 | 22 | import os |
23 | 23 | from datetime import datetime |
24 | 24 |
|
| 25 | +from google.api_core.retry import Retry |
| 26 | + |
25 | 27 | from airflow import models |
26 | 28 | from airflow.providers.google.cloud.operators.dataproc import ( |
27 | 29 | DataprocCreateBatchOperator, |
|
36 | 38 | PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT", "") |
37 | 39 | REGION = "europe-west1" |
38 | 40 | BATCH_ID = f"test-batch-id-{ENV_ID}" |
| 41 | +BATCH_ID_2 = f"test-batch-id-{ENV_ID}-2" |
39 | 42 | BATCH_CONFIG = { |
40 | 43 | "spark_batch": { |
41 | 44 | "jar_file_uris": ["file:///usr/lib/spark/examples/jars/spark-examples.jar"], |
|
58 | 61 | region=REGION, |
59 | 62 | batch=BATCH_CONFIG, |
60 | 63 | batch_id=BATCH_ID, |
61 | | - timeout=5.0, |
| 64 | + ) |
| 65 | + |
| 66 | + create_batch_2 = DataprocCreateBatchOperator( |
| 67 | + task_id="create_batch_2", |
| 68 | + project_id=PROJECT_ID, |
| 69 | + region=REGION, |
| 70 | + batch=BATCH_CONFIG, |
| 71 | + batch_id=BATCH_ID_2, |
| 72 | + result_retry=Retry(maximum=10.0, initial=10.0, multiplier=1.0), |
62 | 73 | ) |
63 | 74 | # [END how_to_cloud_dataproc_create_batch_operator] |
64 | 75 |
|
65 | 76 | # [START how_to_cloud_dataproc_get_batch_operator] |
66 | 77 | get_batch = DataprocGetBatchOperator( |
67 | 78 | task_id="get_batch", project_id=PROJECT_ID, region=REGION, batch_id=BATCH_ID |
68 | 79 | ) |
| 80 | + |
| 81 | + get_batch_2 = DataprocGetBatchOperator( |
| 82 | + task_id="get_batch_2", project_id=PROJECT_ID, region=REGION, batch_id=BATCH_ID_2 |
| 83 | + ) |
69 | 84 | # [END how_to_cloud_dataproc_get_batch_operator] |
70 | 85 |
|
71 | 86 | # [START how_to_cloud_dataproc_list_batches_operator] |
|
80 | 95 | delete_batch = DataprocDeleteBatchOperator( |
81 | 96 | task_id="delete_batch", project_id=PROJECT_ID, region=REGION, batch_id=BATCH_ID |
82 | 97 | ) |
| 98 | + delete_batch.trigger_rule = TriggerRule.ALL_DONE |
| 99 | + |
| 100 | + delete_batch_2 = DataprocDeleteBatchOperator( |
| 101 | + task_id="delete_batch_2", project_id=PROJECT_ID, region=REGION, batch_id=BATCH_ID_2 |
| 102 | + ) |
83 | 103 | # [END how_to_cloud_dataproc_delete_batch_operator] |
84 | 104 | delete_batch.trigger_rule = TriggerRule.ALL_DONE |
85 | 105 |
|
86 | | - create_batch >> get_batch >> list_batches >> delete_batch |
| 106 | + ( |
| 107 | + create_batch |
| 108 | + >> create_batch_2 |
| 109 | + >> get_batch |
| 110 | + >> get_batch_2 |
| 111 | + >> list_batches |
| 112 | + >> delete_batch |
| 113 | + >> delete_batch_2 |
| 114 | + ) |
87 | 115 |
|
88 | 116 | from tests.system.utils.watcher import watcher |
89 | 117 |
|
|
0 commit comments