|
24 | 24 | from airflow.providers.amazon.aws.operators.s3 import S3CreateBucketOperator, S3DeleteBucketOperator |
25 | 25 | from airflow.providers.google.cloud.operators.gcs import GCSCreateBucketOperator, GCSDeleteBucketOperator |
26 | 26 | from airflow.providers.google.cloud.transfers.s3_to_gcs import S3ToGCSOperator |
| 27 | +from airflow.utils.trigger_rule import TriggerRule |
27 | 28 |
|
28 | | -GCP_PROJECT_ID = os.environ.get('GCP_PROJECT_ID', 'gcp-project-id') |
29 | | -S3BUCKET_NAME = os.environ.get('S3BUCKET_NAME', 'example-s3bucket-name') |
30 | | -GCS_BUCKET = os.environ.get('GCP_GCS_BUCKET', 'example-gcsbucket-name') |
31 | | -GCS_BUCKET_URL = f"gs://{GCS_BUCKET}/" |
| 29 | +ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID") |
| 30 | +GCP_PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT") |
| 31 | +DAG_ID = "example_s3_to_gcs" |
| 32 | + |
| 33 | +BUCKET_NAME = f"bucket_{DAG_ID}_{ENV_ID}" |
| 34 | +GCS_BUCKET_URL = f"gs://{BUCKET_NAME}/" |
32 | 35 | UPLOAD_FILE = '/tmp/example-file.txt' |
33 | 36 | PREFIX = 'TESTS' |
34 | 37 |
|
|
37 | 40 | def upload_file(): |
38 | 41 | """A callable to upload file to AWS bucket""" |
39 | 42 | s3_hook = S3Hook() |
40 | | - s3_hook.load_file(filename=UPLOAD_FILE, key=PREFIX, bucket_name=S3BUCKET_NAME) |
| 43 | + s3_hook.load_file(filename=UPLOAD_FILE, key=PREFIX, bucket_name=BUCKET_NAME) |
41 | 44 |
|
42 | 45 |
|
43 | 46 | with models.DAG( |
44 | | - 'example_s3_to_gcs', |
| 47 | + DAG_ID, |
45 | 48 | schedule_interval='@once', |
46 | 49 | start_date=datetime(2021, 1, 1), |
47 | 50 | catchup=False, |
48 | | - tags=['example'], |
| 51 | + tags=['example', 's3'], |
49 | 52 | ) as dag: |
50 | 53 | create_s3_bucket = S3CreateBucketOperator( |
51 | | - task_id="create_s3_bucket", bucket_name=S3BUCKET_NAME, region_name='us-east-1' |
| 54 | + task_id="create_s3_bucket", bucket_name=BUCKET_NAME, region_name='us-east-1' |
52 | 55 | ) |
53 | 56 |
|
54 | 57 | create_gcs_bucket = GCSCreateBucketOperator( |
55 | 58 | task_id="create_bucket", |
56 | | - bucket_name=GCS_BUCKET, |
| 59 | + bucket_name=BUCKET_NAME, |
57 | 60 | project_id=GCP_PROJECT_ID, |
58 | 61 | ) |
59 | 62 | # [START howto_transfer_s3togcs_operator] |
60 | 63 | transfer_to_gcs = S3ToGCSOperator( |
61 | | - task_id='s3_to_gcs_task', bucket=S3BUCKET_NAME, prefix=PREFIX, dest_gcs=GCS_BUCKET_URL |
| 64 | + task_id='s3_to_gcs_task', bucket=BUCKET_NAME, prefix=PREFIX, dest_gcs=GCS_BUCKET_URL |
62 | 65 | ) |
63 | 66 | # [END howto_transfer_s3togcs_operator] |
64 | 67 |
|
65 | 68 | delete_s3_bucket = S3DeleteBucketOperator( |
66 | | - task_id='delete_s3_bucket', bucket_name=S3BUCKET_NAME, force_delete=True |
| 69 | + task_id='delete_s3_bucket', |
| 70 | + bucket_name=BUCKET_NAME, |
| 71 | + force_delete=True, |
| 72 | + trigger_rule=TriggerRule.ALL_DONE, |
67 | 73 | ) |
68 | 74 |
|
69 | | - delete_gcs_bucket = GCSDeleteBucketOperator(task_id='delete_gcs_bucket', bucket_name=GCS_BUCKET) |
| 75 | + delete_gcs_bucket = GCSDeleteBucketOperator( |
| 76 | + task_id='delete_gcs_bucket', bucket_name=BUCKET_NAME, trigger_rule=TriggerRule.ALL_DONE |
| 77 | + ) |
70 | 78 |
|
71 | 79 | ( |
72 | | - create_s3_bucket |
| 80 | + # TEST SETUP |
| 81 | + create_gcs_bucket |
| 82 | + >> create_s3_bucket |
73 | 83 | >> upload_file() |
74 | | - >> create_gcs_bucket |
| 84 | + # TEST BODY |
75 | 85 | >> transfer_to_gcs |
| 86 | + # TEST TEARDOWN |
76 | 87 | >> delete_s3_bucket |
77 | 88 | >> delete_gcs_bucket |
78 | 89 | ) |
| 90 | + |
| 91 | + from tests.system.utils.watcher import watcher |
| 92 | + |
| 93 | + # This test needs watcher in order to properly mark success/failure |
| 94 | + # when "tearDown" task with trigger rule is part of the DAG |
| 95 | + list(dag.tasks) >> watcher() |
| 96 | + |
| 97 | + |
| 98 | +from tests.system.utils import get_test_run # noqa: E402 |
| 99 | + |
| 100 | +# Needed to run the example DAG with pytest (see: tests/system/README.md#run_via_pytest) |
| 101 | +test_run = get_test_run(dag) |
0 commit comments