import logging from tempfile import TemporaryDirectory from airflow.exceptions import AirflowFailException from airflow.providers.amazon.aws.hooks.base_aws import AwsBaseHook from airflow.providers.amazon.aws.hooks.s3 import S3Hook from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook from flows.deezer_daily import api, config from flows.deezer_daily.config import AWS_CONN_ID logger = logging.getLogger(__name__) def fetch(target_s3_path, drop_file_name): """Download a feed file from the drop location into the archive folder. Args: target_s3_path (str): Destination S3 path to the archive location. drop_file_name (str): The dropped file name. """ # setup Zephir settings zephir_settings = dict( username=config.ZEPHIR_CREDENTIALS.get('username'), password=config.ZEPHIR_CREDENTIALS.get('password'), host=config.ZEPHIR_CREDENTIALS.get('host'), path=config.ZEPHIR_CREDENTIALS.get('path') ) with TemporaryDirectory() as temp_dir: local_file_path = api.download_from_zephir( zephir_settings, drop_file_name, temp_dir) s3_hook = S3Hook(aws_conn_id=AWS_CONN_ID) logger.info(f'Uploading file {local_file_path} to S3 {config.ARCHIVE_S3_BUCKET} {target_s3_path}') # noqa s3_hook.load_file( filename=local_file_path, bucket_name=config.ARCHIVE_S3_BUCKET, key=target_s3_path + drop_file_name, replace=True, ) logger.info('Upload done')