From 69343ab054fc5a896eceec5126ac40791b4e223a Mon Sep 17 00:00:00 2001 From: Vincent Beck Date: Mon, 2 May 2022 14:39:50 -0600 Subject: [PATCH 1/2] Add sample dag and doc for S3ListPrefixesOperator --- .../amazon/aws/example_dags/example_s3.py | 39 ++++++++++++------- airflow/providers/amazon/aws/operators/s3.py | 4 ++ .../operators/s3.rst | 16 ++++++++ 3 files changed, 46 insertions(+), 13 deletions(-) diff --git a/airflow/providers/amazon/aws/example_dags/example_s3.py b/airflow/providers/amazon/aws/example_dags/example_s3.py index ecd9d374cf688..1b940c57e9b3f 100644 --- a/airflow/providers/amazon/aws/example_dags/example_s3.py +++ b/airflow/providers/amazon/aws/example_dags/example_s3.py @@ -30,6 +30,7 @@ S3DeleteObjectsOperator, S3FileTransformOperator, S3GetBucketTaggingOperator, + S3ListPrefixesOperator, S3PutBucketTaggingOperator, ) from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor, S3KeysUnchangedSensor @@ -41,6 +42,7 @@ # Empty string prefix refers to the bucket root # See what prefix is here https://docs.aws.amazon.com/AmazonS3/latest/userguide/using-prefixes.html PREFIX = os.environ.get('PREFIX', '') +DELIMITER = os.environ.get('DELIMITER', '/') TAG_KEY = os.environ.get('TAG_KEY', 'test-s3-bucket-tagging-key') TAG_VALUE = os.environ.get('TAG_VALUE', 'test-s3-bucket-tagging-value') DATA = os.environ.get( @@ -106,7 +108,7 @@ def check_fn(files: List) -> bool: # [END howto_operator_s3_delete_bucket_tagging] # [START howto_operator_s3_create_object] - s3_create_object = S3CreateObjectOperator( + create_object = S3CreateObjectOperator( task_id="s3_create_object", s3_bucket=BUCKET_NAME, s3_key=KEY, @@ -115,9 +117,18 @@ def check_fn(files: List) -> bool: ) # [END howto_operator_s3_create_object] + # [START howto_operator_s3_list_prefixes] + list_prefixes = S3ListPrefixesOperator( + task_id="s3_list_prefix_operator", + bucket=BUCKET_NAME, + prefix=PREFIX, + delimiter=DELIMITER, + ) + # [END howto_operator_s3_list_prefixes] + # [START howto_sensor_s3_key_single_key] # Check if a file exists - s3_sensor_one_key = S3KeySensor( + sensor_one_key = S3KeySensor( task_id="s3_sensor_one_key", bucket_name=BUCKET_NAME, bucket_key=KEY, @@ -126,7 +137,7 @@ def check_fn(files: List) -> bool: # [START howto_sensor_s3_key_multiple_keys] # Check if both files exist - s3_sensor_two_keys = S3KeySensor( + sensor_two_keys = S3KeySensor( task_id="s3_sensor_two_keys", bucket_name=BUCKET_NAME, bucket_key=[KEY, KEY_2], @@ -135,7 +146,7 @@ def check_fn(files: List) -> bool: # [START howto_sensor_s3_key_function] # Check if a file exists and match a certain pattern defined in check_fn - s3_sensor_key_function = S3KeySensor( + sensor_key_with_function = S3KeySensor( task_id="s3_sensor_key_function", bucket_name=BUCKET_NAME, bucket_key=KEY, @@ -144,7 +155,7 @@ def check_fn(files: List) -> bool: # [END howto_sensor_s3_key_function] # [START howto_sensor_s3_keys_unchanged] - s3_sensor_keys_unchanged = S3KeysUnchangedSensor( + sensor_keys_unchanged = S3KeysUnchangedSensor( task_id="s3_sensor_one_key_size", bucket_name=BUCKET_NAME_2, prefix=PREFIX, @@ -153,7 +164,7 @@ def check_fn(files: List) -> bool: # [END howto_sensor_s3_keys_unchanged] # [START howto_operator_s3_copy_object] - s3_copy_object = S3CopyObjectOperator( + copy_object = S3CopyObjectOperator( task_id="s3_copy_object", source_bucket_name=BUCKET_NAME, dest_bucket_name=BUCKET_NAME_2, @@ -163,7 +174,7 @@ def check_fn(files: List) -> bool: # [END howto_operator_s3_copy_object] # [START howto_operator_s3_file_transform] - s3_file_transform = S3FileTransformOperator( + transforms_file = S3FileTransformOperator( task_id="s3_file_transform", source_s3_key=f's3://{BUCKET_NAME}/{KEY}', dest_s3_key=f's3://{BUCKET_NAME_2}/{KEY_2}', @@ -174,7 +185,7 @@ def check_fn(files: List) -> bool: # [END howto_operator_s3_file_transform] # [START howto_operator_s3_delete_objects] - s3_delete_objects = S3DeleteObjectsOperator( + delete_objects = S3DeleteObjectsOperator( task_id="s3_delete_objects", bucket=BUCKET_NAME_2, keys=KEY_2, @@ -192,10 +203,12 @@ def check_fn(files: List) -> bool: put_tagging, get_tagging, delete_tagging, - s3_create_object, - [s3_sensor_one_key, s3_sensor_two_keys, s3_sensor_key_function], - s3_copy_object, - s3_sensor_keys_unchanged, - s3_delete_objects, + create_object, + list_prefixes, + [sensor_one_key, sensor_two_keys, sensor_key_with_function], + copy_object, + transforms_file, + sensor_keys_unchanged, + delete_objects, delete_bucket, ) diff --git a/airflow/providers/amazon/aws/operators/s3.py b/airflow/providers/amazon/aws/operators/s3.py index a33913701d7cc..db8b9dc7e119c 100644 --- a/airflow/providers/amazon/aws/operators/s3.py +++ b/airflow/providers/amazon/aws/operators/s3.py @@ -689,6 +689,10 @@ class S3ListPrefixesOperator(BaseOperator): This operator returns a python list with the name of all subfolders which can be used by `xcom` in the downstream task. + .. seealso:: + For more information on how to use this operator, take a look at the guide: + :ref:`howto/operator:S3ListPrefixesOperator` + :param bucket: The S3 bucket where to find the subfolders. (templated) :param prefix: Prefix string to filter the subfolders whose name begin with such prefix. (templated) diff --git a/docs/apache-airflow-providers-amazon/operators/s3.rst b/docs/apache-airflow-providers-amazon/operators/s3.rst index 00f1fe1143353..c856dfeb241b4 100644 --- a/docs/apache-airflow-providers-amazon/operators/s3.rst +++ b/docs/apache-airflow-providers-amazon/operators/s3.rst @@ -196,6 +196,22 @@ To create a new (or replace) Amazon S3 object you can use :start-after: [START howto_operator_s3_create_object] :end-before: [END howto_operator_s3_create_object] +.. _howto/operator:S3ListPrefixesOperator: + +List Amazon S3 prefixes +----------------------- + +To list all Amazon S3 prefixes within an Amazon S3 bucket you can use +:class:`~airflow.providers.amazon.aws.operators.s3.S3ListPrefixesOperator`. +See `here `__ +for more information about Amazon S3 prefixes. + +.. exampleinclude:: /../../airflow/providers/amazon/aws/example_dags/example_s3.py + :language: python + :dedent: 4 + :start-after: [START howto_operator_s3_list_prefixes] + :end-before: [END howto_operator_s3_list_prefixes] + .. _howto/operator:S3CopyObjectOperator: Copy an Amazon S3 object From a1aa1e5b9215e5fc28f748f6598934c5f45aa21e Mon Sep 17 00:00:00 2001 From: Vincent Beck Date: Mon, 9 May 2022 14:46:56 -0600 Subject: [PATCH 2/2] Fix static checks --- airflow/providers/amazon/aws/example_dags/example_s3.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/providers/amazon/aws/example_dags/example_s3.py b/airflow/providers/amazon/aws/example_dags/example_s3.py index 1afb12735d2b1..7e06575d4a58b 100644 --- a/airflow/providers/amazon/aws/example_dags/example_s3.py +++ b/airflow/providers/amazon/aws/example_dags/example_s3.py @@ -30,8 +30,8 @@ S3DeleteObjectsOperator, S3FileTransformOperator, S3GetBucketTaggingOperator, - S3ListPrefixesOperator, S3ListOperator, + S3ListPrefixesOperator, S3PutBucketTaggingOperator, ) from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor, S3KeysUnchangedSensor